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

如何从AWS Glue的Salesforce写入工具获取失败记录?

问题描述

我是AWS Glue新手,正在用Python脚本通过Salesforce写入工具(基于Salesforce连接)把S3数据迁移到Salesforce。预期操作完成后Salesforce会返回jobId/batchId或失败记录,但当前Glue作业返回null。

相关代码

def write_to_salesforce(records, entity_name, id_field_names):
    processed_records = []
    failed_records = []

    if records:
        unique_records = list({record['Name']: record for record in records}.values())
        dynamic_frame = DynamicFrame.fromDF(spark.createDataFrame(unique_records), glueContext, "dynamic_frame")

        logger.info(f"Transformed {entity_name}_dynamic_frame:")
        dynamic_frame.toDF().show()

        try:
            result = glueContext.write_dynamic_frame.from_options(
                frame=dynamic_frame,
                connection_type="salesforce",
                connection_options={
                    "apiVersion": "v60.0",
                    "connectionName": "Salesforce connection",
                    "entityName": entity_name,
                    "writeOperation": "UPSERT",
                    "idFieldNames": id_field_names
                },
                transformation_ctx="dynamic_frame"
            )
            logger.info(f"result: {result}")
            print(f"res: {result.count()}")
            result.toDF().show()
            processed_records.extend(unique_records)
        except Exception as e:
            print(f"Error occurred: {str(e)}")
            logger.error(f"Error writing to {entity_name}: {e}")
            failed_records.extend(e)
    
    return processed_records, failed_records

已尝试的无效方法

  • 带附加参数的作业执行:添加"bulk": true参数以获取批次ID;
  • Scala实现:使用Scala及GlueContext类的getSink或getSinkWithFormat方法写入Salesforce;
  • withErrorsAsDynamicFrame()错误处理:尝试用该方法捕获错误记录为独立动态帧,但Salesforce写入工具不支持此方法;
  • Salesforce调试日志:尝试通过Salesforce调试日志追踪细节,当前使用的是复合API(示例:/services/data/v60.0/composite/sobjects/AddressCatalogRegion__c/ExternalId_t__c);
  • AWS Glue作业书签功能:该功能不支持Salesforce工具;
  • AWS Glue监控与日志:Glue未自动记录Salesforce写入过程的特定错误,仅记录作业状态、转换步骤等基础信息。

请问是否有办法从对应作业中获取失败记录?


解决方案

1. 完善Bulk API配置并查询批次结果

虽然你试过添加"bulk": true,但需要结合Salesforce Bulk API的批次查询逻辑:

  • 在connection_options中明确配置"bulk": true和"batchSize": 1000(根据数据量调整),此时Glue会使用Salesforce Bulk API执行写入,返回的result中会包含批次ID;
  • 拿到批次ID后,调用Salesforce Bulk API的/services/data/v60.0/jobs/ingest/{jobId}/batches/{batchId}/results接口,查询该批次的详细执行结果,结果中会标记每条记录的成功/失败状态及错误原因;
  • 注意:确保Glue作业的IAM角色或Salesforce连接账号拥有Bulk API的访问权限。

2. 自定义对比逻辑捕获失败记录

由于Glue的Salesforce写入工具不支持内置错误帧捕获,可通过"写入后校验"的方式实现:

  • 写入前为每条记录保留唯一标识(比如外部ID或自定义UUID),将原始记录暂存到临时S3路径或Glue临时表;
  • 写入完成后,通过Salesforce连接查询已成功写入的记录的外部ID列表;
  • 对比原始记录列表和查询到的成功列表,未匹配的即为失败记录。示例代码片段:
    import uuid
    
    # 写入前为记录添加临时标识/保留外部ID
    for record in unique_records:
        record['temp_uuid'] = str(uuid.uuid4())
    
    # 写入后查询Salesforce中已存在的外部ID
    success_ids_df = glueContext.create_dynamic_frame.from_options(
        connection_type="salesforce",
        connection_options={
            "apiVersion": "v60.0",
            "connectionName": "Salesforce connection",
            "entityName": entity_name,
            "query": f"SELECT {id_field_names} FROM {entity_name} WHERE {id_field_names} IN ('{','.join([r[id_field_names] for r in unique_records])}')"
        }
    ).toDF()
    
    success_ids = success_ids_df.select(id_field_names).rdd.flatMap(lambda x: x).collect()
    failed_records = [r for r in unique_records if r[id_field_names] not in success_ids]
    

3. 调整日志级别捕获底层API细节

  • 将Glue作业的日志级别调整为DEBUG:在作业配置中修改--log-level参数为DEBUG,或在代码中设置logger.setLevel(logging.DEBUG);
  • 查看CloudWatch日志,其中会打印更多底层Salesforce API的交互细节,包括可能的错误响应内容,帮助定位失败原因。

4. 直接调用Salesforce复合API替代Glue工具

如果Glue的Salesforce工具限制过多,可跳过工具直接用Python调用Salesforce API:

  • 在Glue作业中安装simple-salesforce依赖,批量发送UPSERT请求并直接解析响应中的失败记录:
    from simple_salesforce import Salesforce
    
    # 初始化Salesforce客户端(需从Glue连接中获取实例URL和会话ID)
    sf = Salesforce(instance_url="你的Salesforce实例URL", session_id="你的会话ID")
    
    # 构造复合请求
    composite_request = {
        "allOrNone": False,
        "compositeRequest": [
            {
                "method": "PATCH",
                "url": f"/services/data/v60.0/sobjects/{entity_name}/{id_field_names}/{record[id_field_names]}",
                "body": record
            } for record in unique_records
        ]
    }
    
    # 发送请求并解析失败记录
    response = sf.restful("composite", method="POST", json=composite_request)
    failed_records = [
        composite_request["compositeRequest"][idx]["body"] 
        for idx, res in enumerate(response["compositeResponse"]) 
        if res["httpStatusCode"] >= 400
    ]
    

内容的提问来源于stack exchange,提问作者Poorni_SM

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 12:41:14