Databricks向Cosmos DB批量上传数据过慢的优化方案咨询
优化Cosmos DB数据上传速度的方案
核心问题分析
原代码的低效根源在于:
- 用
collect()将全量数据拉取到驱动节点,既造成内存压力,又完全浪费Spark的分布式处理能力 - 逐条调用
create_items单条插入,未利用Cosmos DB的批量插入能力,产生大量网络开销和请求延迟 - 未使用Spark官方Cosmos DB连接器,无法借助Spark的分布式并行写入特性
具体优化措施
1. 改用Spark Cosmos DB连接器(关键优化)
Spark官方提供的Cosmos DB连接器支持分布式批量写入,直接将DataFrame写入Cosmos DB,彻底避免逐行处理的低效模式。需要先添加对应Spark版本的连接器依赖,例如Spark 3.3版本可使用Maven坐标:com.azure.cosmos.spark:azure-cosmos-spark_3-3_2-12:4.22.0
2. 移除collect(),直接操作DataFrame
collect()会把所有数据拉到单节点处理,完全丧失Spark分布式优势,应直接对DataFrame做字段转换后写入。
3. 配置批量写入参数
通过连接器的配置项优化写入性能,比如调整批次大小、并发数、重试策略等,适配Cosmos DB的处理能力。
4. 直接在DataFrame中生成ID字段
避免逐行解析JSON的开销,直接用Spark的列操作生成id字段。
优化后的代码
import logging from typing import Optional configs = { "dev": { "file_location": "/FileStore/tables/docs/dpidata_pfile_20050523-20221009.csv", "file_type": "csv", "infer_schema": False, "first_row_is_header": True, "delimiter": ",", "cdb_url": "https://xyxxxxxxxxxxxxxxx:443/", "db_name": "abc", "container_name": "dpi", "partition_key": "/dpi", # Cosmos DB连接器配置 "cosmos_write_config": { "spark.cosmos.accountEndpoint": "https://xyxxxxxxxxxxxxxxx:443/", "spark.cosmos.accountKey": "你的Cosmos DB账号密钥", # 需补充实际账号密钥 "spark.cosmos.database": "abc", "spark.cosmos.container": "dpi", # 批量写入优化参数 "spark.cosmos.write.strategy": "Append", "spark.cosmos.write.batch.size": "1000", "spark.cosmos.write.maxConcurrentWrites": "100", "spark.cosmos.write.retry.count": "5", "spark.cosmos.write.retry.delay": "1000" } }, "stg": {}, "prd": {} } class LoadToCdb(): def __init__(self): self.configs = configs["dev"] self.log = logging.getLogger(__name__) logging.basicConfig(level=logging.INFO) def dpi_data_load(self) -> Optional[bool]: try: # 读取CSV文件 df = spark.read.format(self.configs["file_type"]) \ .option("inferSchema", self.configs["infer_schema"]) \ .option("header", self.configs["first_row_is_header"]) \ .option("sep", self.configs["delimiter"]) \ .load(self.configs["file_location"]) # 选择字段、重命名并生成id字段 df = df.select('dpi', 'Entity Type Code') \ .withColumnRenamed("Entity Type Code","entity_type_code") \ .withColumn("id", df["dpi"]) # 写入Cosmos DB df.write.format("cosmos.oltp") \ .options(**self.configs["cosmos_write_config"]) \ .mode("append") \ .save() self.log.info("数据成功写入Cosmos DB") return True except Exception as e: self.log.error("加载CSV到Cosmos DB失败") self.log.error(str(e)) return False if __name__ == "__main__": load_to_cdb = LoadToCdb() load_to_cdb.dpi_data_load()
额外优化建议
- 临时调高Cosmos DB吞吐量:若当前RU(请求单位)不足,可临时提升容器吞吐量,写入完成后再调低,避免限流拖慢速度
- 检查分区键基数:确保分区键
dpi有足够多的不同值,避免热点分区限制写入能力 - 拆分大CSV文件:将9GB的CSV拆分为多个小文件,让Spark的多个分区并行读取处理
- 移除日志打印:原代码中的
print(row)会严重拖慢速度,优化后已移除该操作
内容的提问来源于stack exchange,提问作者GoneCase123
相关产品推荐
相关产品推荐

