如何在Databricks中通过Simple Salesforce捕获批量插入Salesforce的失败记录?
解决方案:直接捕获Salesforce Bulk API返回的失败记录
不需要通过事后关联源数据与Salesforce数据的方式识别失败记录,Simple Salesforce的Bulk API操作本身会返回每条记录的处理结果,直接利用这个返回结果就能高效分离失败记录,完全避免额外查询开销。
实现步骤
执行Bulk插入并获取结果
使用Simple Salesforce的bulk.ObjectName.insert()方法执行数据插入时,该方法会返回一个列表,列表中每个元素对应一条提交记录的处理状态,包含success(布尔值,标记是否成功)、id(成功时的Salesforce记录ID)、errors(失败时的错误详情数组)等字段。关联源记录与处理结果
Salesforce Bulk API保证返回结果的顺序与提交记录的顺序完全一致,因此可以直接将源记录列表与结果列表一一配对,筛选出失败的记录并附加错误信息。保存失败记录
将筛选出的失败记录转换为DataFrame,直接保存到S3指定路径即可。
示例代码
from simple_salesforce import Salesforce from pyspark.sql import SparkSession # 初始化Salesforce连接 sf = Salesforce( username="你的用户名", password="你的密码", security_token="你的安全令牌" ) # 从S3加载源数据 spark = SparkSession.builder.appName("S3ToSalesforce").getOrCreate() source_df = spark.read.csv("s3://你的存储桶/源数据路径/", header=True) # 转换为Simple Salesforce要求的字典列表格式 source_records = source_df.toPandas().to_dict("records") # 执行Bulk插入操作 bulk_response = sf.bulk.Contact.insert( source_records, batch_size=1000, # 根据数据量调整批次大小 use_serial=True ) # 筛选并整理失败记录 failed_records = [] for src_record, result in zip(source_records, bulk_response): if not result["success"]: # 合并错误信息为字符串,方便查看 error_details = "; ".join([err["message"] for err in result["errors"]]) # 保留源数据所有字段,追加错误信息 failed_record = {**src_record, "错误信息": error_details} failed_records.append(failed_record) # 将失败记录保存到S3 if failed_records: failed_df = spark.createDataFrame(failed_records) failed_df.write.mode("overwrite").csv("s3://你的存储桶/失败记录存储路径/", header=True)
关键注意点
- 结果顺序一致性:Salesforce Bulk API严格按照提交记录的顺序返回处理结果,因此
zip配对不会出现错位问题。 - 错误信息提取:
errors字段是数组,每条错误包含message(错误描述)和fields(违规字段),可根据需求选择保存的内容。 - 性能优化:该方式无需额外查询Salesforce,仅利用插入操作的返回结果,不会增加主流程的额外开销。
内容的提问来源于stack exchange,提问作者SK ASIF ALI
相关产品推荐
相关产品推荐

