如何从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
相关产品推荐
相关产品推荐

