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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 02:05:29