Spark3.2.1写入Cosmos DB报Writing job aborted错误
问题背景
使用Spark 3.2.1通过cosmos.oltp连接器往Cosmos DB写入DataFrame数据时,任务失败抛出Writing job aborted异常,驱动日志中根因报错为:
java.lang.IllegalArgumentException: requirement failed: id is a mandatory field. But it is missing or it is not a string
写入代码如下:
df_u.write.format("cosmos.oltp").options(**writeMcgMd).mode("append").save()
对应写入配置:
writeMcgMd = { "spark.cosmos.accountEndpoint" : "https://cccc.azure.com:443/", "spark.cosmos.accountKey" : "ccc", "spark.cosmos.database" : "cccc", "spark.cosmos.container" : "ccc", # "spark.cosmos.write.strategy": "ItemOverwrite" }
从报错打印的待写入数据样例可以看到,所有记录都没有名为id的字段,现有字段为_attachments、databasename、lastloadtime、parentname、query、recordtype、schema、tableName、where_col1、where_col2。
故障原因
Cosmos DB强制要求容器内的每一条文档必须携带字符串类型的id字段作为文档唯一标识,Spark连接器写入前会做前置校验,缺字段、字段类型不是字符串都会直接抛错终止任务。
修复方法
三种方案按需选择即可:
- 写入前手动给DataFrame加字符串类型的
id列,优先用业务字段拼接保证幂等性,没有合适业务字段就生成UUID:
from pyspark.sql import functions as F # 用业务字段拼接id,避免重复写入产生重复文档 df_u = df_u.withColumn("id", F.concat_ws("_", F.col("schema"), F.col("tableName"), F.col("parentname"))) # 无合适业务字段时用自动生成的UUID # df_u = df_u.withColumn("id", F.expr("uuid()")) # 执行写入 df_u.write.format("cosmos.oltp").options(**writeMcgMd).mode("append").save()
- 直接修改写入配置,让连接器自动生成id,不需要手动改DataFrame结构,在配置字典中加一行参数即可:
"spark.cosmos.write.idGenerationStrategy": "Default"
- 如果DataFrame里本身有
id列但不是字符串类型(比如是数字int/long类型),提前把该列转成字符串类型再写入:
df_u = df_u.withColumn("id", F.col("id").cast("string"))
内容的提问来源于stack exchange,提问作者Vaishnavi S
相关产品推荐
相关产品推荐

