You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何高效每日从Teradata向Databricks Spark Delta表导入5000万条数据

问题

在Databricks中搭建ELT流程,每日需从Teradata同步约5000万条单日记录,目标Delta表仅保留一年数据。已优化Teradata查询,EXPLAIN显示查询仅需10分钟,但实际从Teradata读取数据到Spark DataFrame再写入Delta表耗时一整天。由于每次同步的是单日数据,无法用日期列做分区,正考虑其他列作为分区键。

当前代码

数据读取代码

sales_query = f""" SQL Query """

for x in range(0, 10):
  connection_error = False
  print(' Attempt #' + str(x+1))
  try:
    sales = (spark.read.format("jdbc")
        .option("url", database_connection)
        .option("driver", "com.teradata.jdbc.TeraDriver")
        .option("fetchsize","1000000") 
        .option("query", sales_query)
        # .option("partitionColumn","sales_date")  
        # .option("lowerBound", f'{start_date}')
        # .option("upperBound", f'{end_date}')
        # .option("numPartitions", 8)
        .load())
    break
  except Exception as connection_error:
    print('Query Failed. Reason ' + str(connection_error) + '\n' + 'Retrying...')
    sleep(sleep_time)

注:上述读取操作仅耗时20秒,实际仅加载了元数据,未同步全量5000万条记录

数据写入代码

sales.write.format("delta").mode("overwrite").option("overwriteSchema","true").save(f"{project_path}/{level}/{delta_file_name}")

集群配置

  • Databricks版本:12.2 LTS(Apache Spark 3.3.2、Scala 2.12)
  • 节点配置:1个Driver节点(64GB内存、8核);1-2个Worker节点(64-128GB内存、8-16核)
  • 运行时版本:12.2.x-scala2.12

优化建议
  • 启用JDBC并行读取:选择高基数、分布均匀的列(如交易ID、用户ID)作为partitionColumn,设置与Worker总核数匹配的numPartitions(例如2个32核Worker可设为24-32),同时配置对应的lowerBound和upperBound。确保分区列在Teradata上有索引,避免Teradata端查询过载。
  • 调整JDBC fetchsize参数:当前fetchsize=1000000可能过大,易引发内存占用过高问题,建议调整为100000-500000区间,平衡内存占用与网络IO效率。
  • 优化Delta写入逻辑:
    • 移除overwriteSchema=true:除非每日Schema必变,否则该参数会增加额外校验开销,Schema稳定时直接删除即可。
    • 用高频查询列做写入分区:即使是单日数据,也可选择区域ID、交易类型等列作为partitionBy的键,既提升写入并行度,也能优化后续查询的扫描效率。
    • 控制单文件大小:添加.option("maxRecordsPerFile", 1000000)(可根据单条记录大小调整),避免生成大量小文件,提升写入性能。
    • 替换overwrite为merge操作:若目标表仅需新增/更新当日数据,用merge替代全量overwrite,减少不必要的数据重写。
  • 扩容集群资源:
    • 增加Worker节点数量:当前1-2个Worker处理5000万条记录资源不足,建议扩容至3-4个。
    • 调整Executor配置:每个Worker设为16核/64GB(1核对应4GB内存),同时设置spark.executor.instances充分利用节点资源。
    • 开启动态资源分配:让Spark根据任务负载自动调整Executor数量,避免资源浪费或过载。
  • 预加载数据:在写入前执行sales.cache()或sales.count()触发实际数据加载,提前发现读取阶段问题,同时缓存数据避免写入时重复拉取JDBC数据源。

内容的提问来源于stack exchange,提问作者hkay

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 05:40:30