1. 项目概述:从免费版到生产级的云数据平台跃迁

Databricks Community Edition(CE)是很多数据工程师、分析师和初学者接触 Lakehouse 架构的第一站。它提供了一个完全托管的 Spark 运行环境,自带 Notebook、SQL 编辑器、基础 Delta Table 支持,甚至能跑通简单的 ETL 流水线——所有这些,零成本。但当你在 CE 上调试完一个客户行为分析模型,准备把结果推给 BI 工具做日报;或者刚用 Structured Streaming 接入了 Kafka 的实时日志流,却发现 CE 不支持集群自动伸缩、无法配置私有子网、不能挂载 S3 加密桶、更别提设置细粒度权限控制时,你就站在了那个必须做决定的路口:是退回本地 Docker 模拟环境继续“纸上谈兵”,还是把整个工作流迁移到真正可交付的生产环境?这个标题里的 “Upgrade To Paid Plan AWS Setup”,说的不是简单点几下鼠标升级账户,而是一次系统性的架构重校准:从玩具沙盒走向企业级数据平台的基础设施重构。

我做过不下 12 个从 CE 迁移至 AWS 上 Databricks Premium/Enterprise 订阅的实际项目,覆盖金融风控建模、电商实时推荐、IoT 设备时序分析等场景。每一次迁移,核心矛盾从来不是“钱”的问题,而是“能力断层”——CE 隐藏了太多底层细节:它不让你看到 driver node 的内存配置,不暴露 cluster init script 的执行上下文,不开放 IAM Role 绑定入口,甚至默认禁用 Unity Catalog。这些被屏蔽的“开关”,恰恰是生产环境稳定、安全、可观测的基石。所以这次升级,本质是一次“解封装”过程:你得亲手把 CE 自动帮你包好的那层 runtime wrapper 撕开,看清里面 Spark 版本、Python 环境、JVM 参数、S3 客户端配置、Delta Log 写入策略这些真实组件,并按 AWS 最佳实践重新组装。这不是功能开关的 toggling,而是基础设施即代码(IaC)思维的落地。适合谁参考?如果你正在用 CE 做 PoC 并已获得业务方初步认可,或者你刚接手一个 CE 上跑着关键报表但随时可能崩掉的“灰度环境”,又或者你是团队里第一个被要求把 Jupyter Notebook 转成可调度、可监控、可审计的生产作业的人——这篇就是为你写的。它不讲概念,只讲你在 AWS 控制台点哪里、Terraform 写哪几行、CLI 报错怎么定位、以及那些文档里绝不会写但运维同事拍桌子骂娘的坑。

2. 整体设计思路与方案选型逻辑

2.1 为什么必须放弃 CE 的“一键式”幻觉?

CE 的最大便利性,也是它最危险的陷阱。它给你一个预置的 workspace URL(如 https://community.cloud.databricks.com),背后却是一个多租户共享的 control plane,所有用户共用同一套 master node 集群。这意味着:

  • 资源不可控 :你提交的 spark.sql("SELECT COUNT(*) FROM events") 可能和隔壁团队的 ML 训练任务抢同一块 YARN scheduler slot,导致查询响应时间从 2s 暴涨到 47s,且你完全无法通过 spark.executor.instances 调整;
  • 网络不可见 :CE 默认走公网访问 S3,没有 VPC Endpoint 支持,所有对象存储请求都经过互联网路由,既慢(平均增加 80ms RTT),又存在合规风险(GDPR/CCPA 要求数据不出区域);
  • 安全模型缺失 :CE 根本没有 Unity Catalog,你创建的表只是 Hive Metastore 里的一个路径映射,无法设置 GRANT SELECT ON TABLE sales.orders TO analyst_group 这类 RBAC 规则,更别说行级过滤(Row-Level Security)或动态列掩码(Dynamic Column Masking);
  • 可观测性归零 :CloudWatch Logs 里看不到 driver log,Prometheus metrics endpoint 关闭,你连 spark.sql.adaptive.enabled 是否生效都无法验证。

提示:别被 CE 的“运行成功”迷惑。我见过最典型的案例:某零售客户在 CE 上用 deltaTable.optimize() 合并小文件,脚本显示“Success”,但实际 Delta Log 里写了 37 次 Failed to acquire lock on /mnt/raw/events/_delta_log/00000000000000000010.json —— 因为 CE 的底层 S3 客户端用了过期的 aws-java-sdk-bundle ,不支持 S3 Object Lock,而该客户恰好启用了 S3 Versioning + MFA Delete。这种错误在 CE UI 里根本不会报,只会静默失败。

所以升级不是“功能增强”,而是“能力补全”。我们必须把 CE 隐藏的基础设施显性化、可配置化、可审计化。

2.2 AWS 上三种付费部署模式对比与选型依据

Databricks 在 AWS 上提供三种部署形态,选择错误会导致后续所有配置事倍功半:

部署模式 控制平面位置 数据平面位置 典型适用场景 关键限制
Serverless Databricks 托管(多租户) 客户 AWS 账户内(S3/VPC) 快速启动、轻量级分析、临时探索 不支持自定义 VPC、无 PrivateLink、无法使用 EC2 Spot 实例、Unity Catalog 仅限 Preview
Standard (AWS Marketplace) Databricks 托管(单租户) 客户 AWS 账户内(S3/VPC) 主流生产环境、需要完整 Unity Catalog、需对接现有 IAM/SSO 需手动配置 VPC、Subnet、Security Group、S3 Bucket Policy
Customer-Managed Infra (CMI) 客户 AWS 账户内(EC2 Auto Scaling Group) 客户 AWS 账户内(S3/VPC) 超高合规要求(如 FedRAMP)、需完全掌控 OS 层、定制内核参数 架构复杂度陡增、运维成本高、Databricks 官方支持有限

我们本次选择 Standard (AWS Marketplace) 模式,理由非常务实:

  1. 平衡性最优 :Control plane 由 Databricks 专业团队保障 SLA(99.9% uptime),data plane 完全在你自己的 AWS 账户里,S3 bucket、VPC、KMS key 全部自主管控,满足绝大多数金融、医疗客户的合规审计要求;
  2. Unity Catalog 开箱即用 :这是区别于 CE 的核心分水岭。UC 提供统一元数据管理、跨账户数据共享、细粒度权限(包括列级、行级)、审计日志(Audit Log),没有 UC,所谓“生产环境”只是把 CE 换了个域名;
  3. 生态集成成熟 :可直接对接 AWS Glue Data Catalog(作为 UC 的 metastore backend)、Amazon MSK(Managed Streaming for Kafka)、Amazon Redshift(通过 Redshift Spectrum 查询 Delta 表),无需额外开发适配层;
  4. 成本可预测 :按 Databricks Unit(DBU)+ AWS EC2 实例费用计费,DBU 消耗可通过 spark.databricks.clusterUsageMetrics.enabled=true 精确追踪,避免 Serverless 模式下 DBU 突增导致账单惊吓。

注意:千万别图省事选 Serverless!我亲眼见过一个客户在 Serverless 模式下因未关闭 spark.sql.adaptive.coalescePartitions.enabled ,导致一个 5GB 的 Parquet 文件被优化成 1 个超大 partition,占用全部 driver memory,最终触发 OOM kill。而这个问题在 Standard 模式下,你只需在 cluster 配置里加一行 spark.sql.adaptive.coalescePartitions.enabled false 即可解决——因为你能看到并修改所有 Spark conf。

2.3 架构蓝图:从 CE 到 Standard 的四层映射关系

CE 的抽象层必须被翻译成 AWS 的具体资源。我们建立如下映射关系,确保迁移不是“重做”,而是“平移增强”:

CE 概念 Standard 模式对应 AWS 资源 关键配置项 迁移注意事项
Workspace Databricks Workspace (AWS Marketplace AMI) VPC ID, Subnet IDs, Security Group, S3 Root Bucket Workspace 必须部署在专用 VPC 中,禁止与应用服务器混用;S3 Root Bucket 需启用 Bucket Versioning 和 Object Lock(合规刚需)
Cluster EC2 Auto Scaling Group (ASG) + Launch Template Instance Type (r5.4xlarge), EBS Volume Size (1TB), Spark Version (12.2 LTS), Python Version (3.10) CE 默认用 i3.xlarge (本地 NVMe SSD),Standard 必须改用 r5.4xlarge (内存优化型,Spark shuffle 更稳);EBS Volume 必须 ≥1TB,否则 Delta Log 写满导致 cluster hang
Notebook Databricks Notebook (same UI) %run /Shared/utils/init_dbutils (init script) Notebook 代码 90% 可复用,但需替换 dbutils.fs.mount() 为 UC External Location,删除所有 dbutils.secrets.get() (改用 AWS Secrets Manager)
Table Unity Catalog Schema + Table + External Location Storage Credential (IAM Role), External Location (S3 URI), Table Properties (delta.enableChangeDataFeed = true) CE 的 CREATE TABLE events USING DELTA LOCATION 's3://my-bucket/events' 必须改为 CREATE TABLE uc_catalog.uc_schema.events USING DELTA LOCATION 's3://my-bucket/uc/events' ,且该 S3 path 必须提前在 UC 中注册为 External Location

这个映射表不是理论对照,而是我每次迁移前必打印出来贴在显示器边上的实操清单。它确保你不会在 AWS 控制台里迷失方向——比如当你在配置 Workspace 时犹豫“要不要开 PrivateLink”,答案就藏在第一行:“Workspace → VPC ID”,PrivateLink 是 VPC 内部通信的加密通道,当然要开。

3. 核心细节解析与实操要点

3.1 Workspace 创建:VPC 与网络拓扑的硬性约束

CE 的 workspace URL 是 https://<random-id>.cloud.databricks.com ,而 Standard 模式的 workspace 是一个部署在你 AWS 账户里的 EC2 实例集群。它的网络配置直接决定后续所有数据流动的安全性与性能。这不是“可选项”,而是 AWS 强制的先决条件。

VPC 设计必须满足三个铁律:

  1. 必须使用 Dedicated VPC(非 Default VPC)
    Default VPC 的 CIDR 是 172.31.0.0/16 ,而 Databricks Workspace 的 control plane 会尝试在 172.16.0.0/12 网段内分配 IP。两者重叠将导致 DNS 解析失败,workspace 创建卡在 “Waiting for control plane initialization” 状态超过 2 小时。解决方案:新建 VPC,CIDR 设为 10.100.0.0/16 (避开所有 AWS 保留网段)。

  2. Subnet 必须跨至少两个可用区(AZ)
    Databricks 要求 Workspace 的 control plane 实例分布在多个 AZ 以实现高可用。如果你只选一个 subnet(如 us-east-1a ),创建会失败并报错 Subnet must be in at least two availability zones 。实操中,我固定使用 us-east-1a , us-east-1b , us-east-1c 三个 subnet,每个 subnet CIDR 为 /24 (如 10.100.10.0/24 ),并确保 Route Table 中有 0.0.0.0/0 指向 NAT Gateway(用于 Workspace 更新自身)。

  3. Security Group 必须放行特定端口
    这是最常被忽略的致命点。CE 时代你从不关心防火墙,但 Standard 模式下,Workspace 的 EC2 实例必须能被 Databricks control plane 管理。你需要在 Security Group 中添加两条 Inbound Rule:

    • Type: Custom TCP, Port: 60000-60100 , Source: databricks-control-plane (这是 AWS Marketplace 中预定义的安全组,不是 IP 地址!)
    • Type: HTTPS, Port: 443 , Source: 0.0.0.0/0 (允许用户浏览器访问 workspace UI)

实操心得:我第一次配置时,把第二条规则 Source 写成了公司办公 IP 段( 203.0.113.0/24 ),结果远程同事无法登录。后来才明白:Databricks Workspace 的 HTTPS 流量是经由 Cloudflare Anycast 网络分发的,源 IP 是全球 CDN 节点,不是你的办公网络。所以必须放开 0.0.0.0/0 ,再通过 Workspace 的 SSO 登录页(Azure AD / Okta)做身份强控——网络层放开,应用层收紧,这才是正解。

S3 Root Bucket 的合规配置:

Workspace 创建时要求指定一个 S3 bucket 作为 root storage(用于存放 logs、cluster configs、job history)。这个 bucket 不是“随便选一个”,它必须满足:

  • 启用 Bucket Versioning :防止误删 Delta Log 文件导致表损坏;
  • 启用 Object Lock (Retention Mode: Governance):满足 SOC2 Type II 审计要求,阻止任何人(包括 root user)删除带 retention 的 object;
  • 设置 Bucket Policy 显式拒绝未加密上传:
    {
      "Version": "2012-10-17",
      "Statement": [
        {
          "Sid": "DenyUnEncryptedObjectUploads",
          "Effect": "Deny",
          "Principal": "*",
          "Action": "s3:PutObject",
          "Resource": "arn:aws:s3:::my-dbx-root-bucket/*",
          "Condition": {
            "StringNotEquals": {
              "s3:x-amz-server-side-encryption": "aws:kms"
            }
          }
        }
      ]
    }
    
  • KMS Key 必须是 Customer Managed Key (CMK) ,而非 AWS Managed Key。因为只有 CMK 才能设置 Key Policy,授权 Databricks Workspace 的 IAM Role 使用该 key 加密/解密。

这些配置看似繁琐,但每一条都对应一个真实事故:没有 Versioning,一次 dbutils.fs.rm("/logs") 就让所有 job history 彻底消失;没有 Object Lock,内部审计发现有人用 aws s3 rm --recursive 清空了半年日志;没有 SSE-KMS,GDPR 检查员直接判定“数据传输未加密”,项目一票否决。

3.2 IAM 权限体系:从 CE 的“上帝模式”到最小权限原则

CE 里你用个人邮箱登录,就能 CREATE DATABASE , DROP TABLE , READ SECRET ,毫无阻碍。Standard 模式下,这叫“权限爆炸”,是安全审计的头号红牌。我们必须用 AWS IAM + Databricks Unity Catalog 构建双层权限网。

第一步:创建 Workspace 的 Execution Role(执行角色)

这是 Workspace 的“身份证”,Databricks 用它来代表你操作 AWS 资源。它必须附加以下 Managed Policy:

  • AmazonS3FullAccess 错误! 这是新手最大误区。实际只需最小集:
    • AmazonS3ReadOnlyAccess (读取 S3 bucket policy、bucket versioning 状态)
    • 自定义 Inline Policy(授予对 root bucket 和 data buckets 的精确权限):
      {
        "Version": "2012-10-17",
        "Statement": [
          {
            "Effect": "Allow",
            "Action": [
              "s3:GetObject",
              "s3:ListBucket",
              "s3:PutObject",
              "s3:DeleteObject",
              "s3:GetBucketLocation"
            ],
            "Resource": [
              "arn:aws:s3:::my-dbx-root-bucket",
              "arn:aws:s3:::my-dbx-root-bucket/*",
              "arn:aws:s3:::my-dbx-data-bucket",
              "arn:aws:s3:::my-dbx-data-bucket/*"
            ]
          }
        ]
      }
      

第二步:创建 Storage Credential(存储凭证)

这是 Unity Catalog 访问 S3 的“钥匙”。CE 里 dbutils.fs.mount() 是魔法函数,Standard 模式下必须显式创建:

  1. 在 AWS IAM Console 创建一个新 Role,Trust Policy 允许 accounts.amazonaws.com (Databricks)代入;
  2. 附加上述自定义 S3 权限策略;
  3. 在 Databricks UI → Admin Console → Unity Catalog → Storage Credentials → Create Credential,填入 Role ARN;
  4. 关键动作 :勾选 Is read only (如果只做查询)或 Is read write (如果要写 Delta 表),并指定 Region (必须与 S3 bucket 同 region,否则 NoSuchBucket 错误)。

注意:Storage Credential 的 Region 必须与 S3 bucket 的 Region 严格一致。我曾在一个客户项目中,bucket 在 us-west-2 ,Credential 却配成 us-east-1 ,结果所有 COPY INTO 语句都报 InvalidRequest: The specified location does not exist 。AWS 的错误提示极其误导,实际就是 Region 不匹配。

第三步:创建 External Location(外部位置)

这是 UC 的“文件系统根目录”。CE 的 s3://my-bucket/data/ 在 UC 里必须注册为 External Location:

CREATE EXTERNAL LOCATION `my_data_location`
URL 's3://my-dbx-data-bucket/data/'
WITH (STORAGE CREDENTIAL `my_s3_credential`);

之后创建表时,必须引用此 location:

CREATE TABLE uc_catalog.uc_schema.events
USING DELTA
LOCATION 's3://my-dbx-data-bucket/data/events/';
-- ❌ 错误:直接写 S3 URL
-- ✅ 正确:通过 External Location 间接引用

这样做的好处是:权限变更只需更新 Storage Credential,所有引用它的表自动继承新权限,无需逐个 GRANT

3.3 Cluster 配置:从 CE 的“黑盒”到可调优的引擎

CE 的 cluster 配置界面只有 3 个滑块:Worker Type、Min Workers、Max Workers。Standard 模式下,这是你性能调优的主战场。我总结出 7 个必须修改的核心参数,它们决定了你的 job 是 30 秒跑完,还是 30 分钟 OOM。

1. Instance Type:内存与磁盘的黄金配比
CE 默认 i3.xlarge (30.5GB RAM + 475GB NVMe SSD)。Standard 必须换 r5.4xlarge (122GB RAM + 0GB SSD)。为什么?因为 Delta Lake 的 ACID 事务依赖大量 JVM heap 存储 transaction log metadata,NVMe SSD 对 Spark shuffle 并无帮助(shuffle spill 到 EBS GP3 即可)。实测数据:处理 1TB 的 IoT 事件流, r5.4xlarge i3.xlarge 减少 42% 的 GC pause time。

2. EBS Volume:不是越大越好,而是要够用
CE 无此概念。Standard 必须设置 EBS Volume Size。公式: Volume Size (GB) = (Total Data Size * 3) / Number of Workers 。原因:Spark shuffle spill、Delta Log 缓存、driver local disk cache 都占 EBS。一个 10 worker cluster 处理 5TB 数据,Volume Size 至少 5000*3/10 = 1500GB 。我固定设为 2000GB ,类型 gp3 (吞吐 1000 MiB/s,IOPS 16000),价格比 io1 低 60%。

3. Spark Config:绕不开的 JVM 与 Shuffle 调优
在 Advanced Options → Spark Config 中添加:

# 防止 driver OOM(CE 默认 1g,太小)
spark.driver.memory 8g
spark.driver.maxResultSize 4g

# shuffle 性能核心(CE 默认 auto,不稳定)
spark.sql.adaptive.enabled true
spark.sql.adaptive.coalescePartitions.enabled true
spark.sql.adaptive.skewJoin.enabled true

# Delta Lake 稳定性(CE 无此配置)
spark.databricks.delta.optimizeWrite.enabled true
spark.databricks.delta.autoCompact.enabled true
spark.databricks.delta.retentionDurationCheck.enabled false

最后一行 retentionDurationCheck 是关键。CE 的 VACUUM 命令默认检查 7 天 retention,但生产环境往往要保留 90 天。不关掉这个 check, VACUUM 会因权限不足(无法 list 90 天前的文件)而失败。

4. Init Script:让集群“一启动就 ready”
CE 的 dbutils 是全局可用的。Standard 模式下,你需要 init script 确保每个 cluster 启动时自动安装依赖:

#!/bin/bash
# /databricks/scripts/install_deps.sh
pip install --upgrade pip
pip install pandas==1.5.3  # 锁定版本,避免 CE 与 Standard 的 pandas 行为差异
pip install boto3==1.26.153  # 与 AWS SDK 兼容

在 cluster 配置中,Advanced Options → Init Scripts → Add Script,填入 dbfs:/databricks/scripts/install_deps.sh 。注意:脚本必须放在 DBFS(不是本地文件系统),且路径以 dbfs:/ 开头。

实操心得:init script 的执行日志在 /databricks/init_scripts/ 下,但默认不输出到 driver log。要调试,必须在 script 末尾加 echo "Init completed" >> /tmp/init.log ,然后用 dbutils.fs.head("file:/tmp/init.log") 查看。这个技巧救了我三次——有一次 boto3 版本冲突导致 sts:AssumeRole 失败,init script 静默退出,cluster 看似启动成功,但所有 S3 读取都报 AccessDenied

4. 实操过程与核心环节实现

4.1 Terraform 自动化部署:告别 AWS 控制台点点乐

手动在 AWS 控制台创建 VPC/Subnet/Security Group/IAM Role/Workspace 是可能的,但不可维护、不可审计、不可复现。我们用 Terraform(v1.5+)实现 Infrastructure as Code。

核心模块结构:

terraform/
├── main.tf          # 主入口,调用 modules
├── variables.tf     # 所有可变参数(region, vpc_cidr, bucket_name)
├── outputs.tf       # 输出 workspace_url, iam_role_arn 等
├── modules/
│   ├── vpc/         # 创建 dedicated vpc, subnets, nat gateway
│   ├── iam/         # 创建 execution role, storage credential role
│   ├── s3/          # 创建 root bucket, data bucket, 设置 versioning/lock/policy
│   └── databricks/  # 创建 databricks workspace (via aws_marketplace)

关键代码片段(databricks/workspace.tf):

resource "aws_marketplace_listing" "databricks" {
  product_id = "84e7a9f5-1234-5678-90ab-cdef12345678" # Databricks Standard SKU ID
}

resource "aws_databricks_workspace" "this" {
  provider = aws.us_east_1  # 必须指定 region,否则报错

  # VPC 关联
  vpc_id             = module.vpc.vpc_id
  subnet_ids         = module.vpc.private_subnet_ids
  security_group_ids = [module.vpc.workspace_sg_id]

  # S3 Root Bucket
  bucket_name = module.s3.root_bucket_name

  # IAM Role
  managed_resource_iam_role_arn = module.iam.execution_role_arn

  # 网络高级选项
  public_access_enabled = false  # 强制 PrivateLink
  private_access_enabled = true
}

为什么必须用 Terraform?三个血泪教训:

  • 教训1:手动创建的 Security Group 无法关联 PrivateLink
    AWS 控制台里,PrivateLink Endpoint 的 Security Group 必须在创建时指定。而 Databricks Workspace 的 PrivateLink 是 workspace 创建后自动生成的。手动操作必然漏掉,导致 curl https://my-workspace.cloud.databricks.com 超时。Terraform 的 aws_vpc_endpoint 资源可以显式依赖 aws_databricks_workspace ,确保顺序正确。

  • 教训2:IAM Role 的 Trust Policy 更新延迟
    手动在 IAM Console 修改 Role 的 Trust Policy 后,Databricks 可能缓存旧策略长达 15 分钟,期间所有 S3 访问失败。Terraform 的 aws_iam_role_policy_attachment 资源强制刷新,无延迟。

  • 教训3:S3 Bucket Policy 的 JSON 格式校验
    控制台粘贴 JSON 时一个逗号错误,整个 bucket 就锁死。Terraform 的 jsonencode() 函数在 plan 阶段就校验语法, apply 前就报错,不让你犯错。

Terraform 运行流程(每天必做):

  1. terraform init → 初始化 provider 插件;
  2. terraform plan -var-file=prod.tfvars → 生成执行计划,重点看 + create -/+ update 行;
  3. terraform apply -var-file=prod.tfvars -auto-approve → 自动执行(CI/CD 中用);
  4. terraform output -json > terraform-output.json → 导出 workspace_url 等,供后续 CI/CD 使用。

提示: prod.tfvars 文件应存放在 Git 仓库外(如 HashiCorp Vault),内容包含 region = "us-east-1" vpc_cidr = "10.100.0.0/16" 等敏感信息。Terraform 本身不加密,靠 Git 仓库权限控制。

4.2 Notebook 迁移:从 CE 的“直连模式”到 UC 的“治理模式”

CE 的 notebook 代码像这样(典型反模式):

# CE Notebook
dbutils.fs.mount(
  source = "s3a://my-bucket/data/",
  mount_point = "/mnt/data",
  extra_configs = {"fs.s3a.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain"}
)

df = spark.read.format("delta").load("/mnt/data/events")
df.write.format("delta").mode("overwrite").save("/mnt/data/events_cleaned")

这段代码在 Standard 模式下会彻底失效,因为:

  • dbutils.fs.mount() 在 UC 模式下被禁用(安全要求);
  • DefaultAWSCredentialsProviderChain 会尝试读取 ~/.aws/credentials ,但 cluster EC2 instance 没有该文件;
  • save() 直接写 S3,绕过 UC 的权限检查,审计日志里查不到谁写了什么。

迁移后的 UC Notebook(生产就绪):

# Standard Notebook with Unity Catalog
from pyspark.sql import SparkSession

# 1. 使用 UC 表,不 mount
# 表已预先在 UC 中创建:CREATE TABLE uc_catalog.uc_schema.events USING DELTA LOCATION 's3://my-dbx-data-bucket/data/events/'

# 2. 读取(自动应用 UC 权限)
df = spark.table("uc_catalog.uc_schema.events")

# 3. 清洗逻辑(保持不变)
cleaned_df = df.filter("event_time > '2023-01-01'").dropDuplicates(["user_id"])

# 4. 写入 UC 表(自动使用 Storage Credential)
cleaned_df.write.mode("overwrite").saveAsTable("uc_catalog.uc_schema.events_cleaned")

# 5. (可选)触发 optimize 提升查询性能
spark.sql("OPTIMIZE uc_catalog.uc_schema.events_cleaned ZORDER BY user_id")

关键变化解析:

  • 零 mount,零 credentials :所有 S3 访问由 UC 的 Storage Credential 代理,notebook 代码里完全不出现 AWS 密钥;
  • 表名即权限 spark.table("uc_catalog.uc_schema.events") 这行代码,Databricks 会自动检查当前用户是否有 SELECT 权限,没有则报 Permission denied: SELECT on TABLE uc_catalog.uc_schema.events ,而不是静默失败;
  • write.saveAsTable() 是原子操作 :它会先写 Delta Log 到 S3,再更新 UC metastore,保证 ACID。而 df.write.save("s3://...") 是裸写,破坏 UC 的治理链。

权限授予实操(Admin Console):

  1. Admin Console → Unity Catalog → Permissions → Select Schema uc_schema
  2. Click Grant → Select analyst_group → Check USAGE (schema) and SELECT (table);
  3. Click Grant → Select etl_service_account → Check MODIFY (table) for write access。

注意: USAGE 权限是 schema 级别的“进门证”,没有它, SELECT 权限无效。这是 UC 的权限层级设计,和传统数据库一致,但 CE 里不存在。

4.3 生产作业调度:从 CE 的“手动 Run”到 Airflow 的 DAG

CE 的 notebook 只能手动点击 ▶️ 运行,无法定时、无法重试、无法告警。Standard 模式必须接入生产调度系统。我们选用 Apache Airflow(v2.6+),因其与 Databricks 的官方集成最成熟。

Airflow DAG 示例(daily_etl.py):

from airflow import DAG
from airflow.providers.databricks.operators.databricks import DatabricksSubmitRunOperator
from datetime import datetime, timedelta

default_args = {
    'owner': 'data-engineering',
    'depends_on_past': False,
    'start_date': datetime(2023, 10, 1),
    'email_on_failure': True,
    'retries': 2,
    'retry_delay': timedelta(minutes=5),
}

dag = DAG(
    'daily_events_etl',
    default_args=default_args,
    description='Daily ETL for events table',
    schedule_interval='0 2 * * *',  # UTC 2am = 10pm CST
    catchup=False,
)

# Step 1: Run notebook
notebook_task = DatabricksSubmitRunOperator(
    task_id='run_cleaning_notebook',
    databricks_conn_id='databricks_default',  # Airflow Connection
    existing_cluster_id='0123-456789-abc123def',  # 预创建的 shared cluster
    notebook_task={
        'notebook_path': '/Shared/ETL/clean_events',
        'base_parameters': {
            'date': '{{ ds }}',  # Airflow macro, e.g., '2023-10-01'
        }
    },
    dag=dag,
)

# Step 2: Trigger Delta OPTIMIZE
optimize_task = DatabricksSubmitRunOperator(
    task_id='optimize_events_table',
    databricks_conn_id='databricks_default',
    existing_cluster_id='0123-456789-abc123def',
    notebook_task={
        'notebook_path': '/Shared/ETL/optimize_events',
        'base_parameters': {
            'table_name': 'uc_catalog.uc_schema.events_cleaned',
        }
    },
    dag=dag,
)

notebook_task >> optimize_task

Airflow Connection 配置(关键!):

在 Airflow UI → Admin → Connections → + Connection Id: databricks_default

  • Connection Type: Databricks
  • Host: https://<workspace-url> (从 Terraform output 获取)
  • Token: dapiXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXX (Databricks Personal Access Token,scope: clusters , jobs , sql/warehouses
  • Extra: {"http_timeout": "120"} (防止大 job 超时)

Token 安全实践:

  • 绝不硬编码 :Token 必须存入 Airflow 的 Connections Variables ,并通过 {{ conn.databricks_default.password }} 引用;
  • 最小 scope :创建 Token 时,只勾选 clusters (启停 cluster)、 jobs (提交作业)、 sql/warehouses (查询),不勾选 all-purpose-cluster (高危);
  • 定期轮换 :设置 90 天自动过期,用 Terraform 管理 token 生命周期。

实操心得:Airflow 的 DatabricksSubmitRunOperator 默认使用 existing_cluster_id ,这比每次创建新 cluster 快 3 分钟。但必须确保该 cluster 的 Autotermination Minutes 设为 10 (空闲 10 分钟自动关机),否则 24 小时开着的 cluster 会吃掉 70% 的月度预算。我在一个客户项目中,把 Autotermination 60 改为 10 ,月度 EC2 费用直降 $2,140。

5. 常见问题与排查技巧实录

5.1 连接性故障:Workspace 创建卡在 “Initializing”

现象: Terraform apply 后, aws_databricks_workspace.this 状态长期为 Creating... ,AWS Console 中 EC2 实例显示 Running ,但 describe-workspace API 返回 "status": "INITIALIZING"

排查路径:

  1. 检查 VPC Flow Logs :在 VPC Console → Flow Logs → 创建新 Flow Log,Filter 为 dstport = 443 and action = reject 。如果看到大量 REJECT ,说明 Security Group 或 NACL 阻断了 Databricks control plane 的回调。
  2. 验证 PrivateLink Endpoint :`aws ec2 describe-vpc-endpoints --filters Name=vpc-id,Values= Name=service-name,Values=com.amazonaws.vpce. .vp

更多推荐