Snowflake到Databricks并行数据下载性能问题及优化咨询
问题描述
我在Snowflake中有一张包含10B条记录的大表,希望通过Snowflake Connector(spark.read.format("snowflake"))将其导入Databricks。尝试通过日期列拆分表实现并行获取,并使用Databricks的并发Notebook机制来运行任务,拆分代码如下:
val notebooks = Seq( NotebookData("my_snowflake_table", 6000,Map("start_date" -> "2022-05-01", "end_date" -> "2022-08-01")), NotebookData("my_snowflake_table", 6000,Map("start_date" -> "2022-08-01", "end_date" -> "2022-11-01")), NotebookData("my_snowflake_table", 6000,Map("start_date" -> "2022-11-01", "end_date" -> "2022-12-01"))) // Run the notebooks in parallel val res = parallelNotebooks(notebooks) Await.result(res, 7200 seconds) // this is a blocking call. res.value
原本期望通过自动扩集群实现水平扩展,但实际并未达到预期:单独运行每个拆分任务约需15分钟,但并行运行时1小时都无法完成。
疑问:
- 是否因为所有拆分任务都在同一个Driver节点运行,导致Driver节点带宽成为瓶颈,进而影响整体性能?
- 该设计是否符合Databricks与Snowflake的协作逻辑?
- 还有哪些更快的方法可以将Snowflake数据导入Databricks?
问题解答
1. Driver节点带宽是否是瓶颈?
是的,这种多Notebook并发的方式大概率会让Driver成为性能瓶颈。每个Notebook任务若由同一个Driver发起Snowflake连接、处理元数据协调和数据传输,当多个任务并行时,Driver的网络带宽、CPU、内存资源会被快速耗尽,无法高效处理多并发的数据拉取请求,直接拖慢整体执行速度。
2. 该设计是否符合Databricks与Snowflake的协作逻辑?
不符合。Databricks与Snowflake的原生协作逻辑是利用Spark的分布式读取能力,而非通过多Notebook并发模拟并行。Spark本身是分布式计算框架,应该让Executor节点直接与Snowflake建立连接并行拉取数据,手动拆分Notebook的方式反而破坏了Spark的分布式优势,增加了Driver层的不必要开销。
3. 更快的导入方法
(1)利用Snowflake Connector原生并行读取能力
Snowflake的Spark Connector支持通过分区参数自动实现分布式读取,让Executor直接并行拉取数据,无需手动拆分Notebook:
val sfOptions = Map( "sfUrl" -> "your-snowflake-url", "sfUser" -> "user", "sfPassword" -> "password", "sfDatabase" -> "db", "sfSchema" -> "schema", "sfWarehouse" -> "warehouse", "dbtable" -> "my_snowflake_table", "partition_column" -> "date_col", // 你的日期列 "lower_bound" -> "2022-05-01", "upper_bound" -> "2022-12-01", "num_partitions" -> 12 // 根据数据量和集群规模调整 ) val df = spark.read.format("snowflake").options(sfOptions).load() df.write.format("delta").save("/databricks/delta/my_table")
这种方式下,Spark会自动将数据拆分为多个分区,每个Executor负责一个分区的数据拉取,充分利用集群分布式资源,避免Driver瓶颈。
(2)Snowflake Unload到云存储,再导入Databricks
对于10B+记录的超大规模表,更高效的方式是先将Snowflake数据Unload到云存储(S3/ADLS/GCS)的Parquet/ORC格式,再由Databricks直接读取:
- 在Snowflake中执行Unload命令:
COPY INTO '@your_stage/path/to/data' FROM my_snowflake_table FILE_FORMAT = (TYPE = PARQUET COMPRESSION = SNAPPY) PARTITION BY (date_col) OVERWRITE = TRUE;
- 在Databricks中读取云存储数据:
val df = spark.read.parquet("s3://your-bucket/path/to/data") // 对应云存储路径 df.write.format("delta").saveAsTable("my_table")
Snowflake的Unload是并行化执行的,且云存储读取性能远高于跨平台直接拉取,还能利用Databricks的云原生读取优化。
(3)调整资源配置优化
- 扩容Snowflake Warehouse,确保其能支撑并行读取请求,避免Snowflake侧成为瓶颈;
- 配置Databricks集群为足够的Executor资源,选择高网络带宽的实例类型;
- 开启Snowflake Connector的批量读取优化,比如设置
batch_size参数增大单次拉取的数据量。
内容的提问来源于stack exchange,提问作者Abhishek

