Airflow中PlainXComArg传入SalesforceBulkOperator前未解析引发TypeError问题求助
Airflow中PlainXComArg传入SalesforceBulkOperator前未解析引发TypeError问题求助
我现在在写一个把附件上传到Salesforce的Airflow流水线,遇到了个头疼的问题。我用Taskflow装饰的函数生成了要给SalesforceBulkOperator用的映射关系,本来以为直接把返回的变量传过去就行,结果跑起来就报错:TypeError: object of type 'PlainXComArg' has no len(),看起来Airflow在把参数传给Operator之前没有解析XComArg。
我的代码一开始是这样的:
mappings = generate_mappings(argument) #@task decorated function bulk_insert = SalesforceBulkOperator( salesforce_conn_id='salesforce_conn_id', task_id='bulk_insert', operation='insert', object_name='Attachment', payload=mappings, external_id_field='Id', batch_size=10000, use_serial=False, ) mappings >> bulk_insert
我查了不少文档(比如Astronomer和Airflow官方的Taskflow与传统Operator配合的示例),按照文档里的写法应该能正常工作,但实际就是不行。
我知道如果把这个Operator放进一个Taskflow函数里,手动调用.execute(context)也能让代码跑起来,但总觉得这种写法不够干净利落,不想这么做。
到底是我哪里写错了?还是我误解了文档的意思?难道这种传参方式只对部分Operator生效吗?
另外有个细节:我必须手动显式声明依赖关系,否则Airflow不知道要等generate_mappings执行完再调用Operator,这会不会是个线索?
编辑:应要求附上完整的相关代码:
@task def generate_mappings(tuples): mappings = [] for tuple in tuples: mapping = { 'id': None, 'ParentID': tuple[0], 'Name': attachment_name, 'ContentType': 'application/pdf', 'Body': tuple[1], } mappings.append(mapping) return mappings @dag( schedule=None, start_date=pendulum.datetime(1970, 1, 1, tz='UTC'), catchup=False, ) def dag(): mappings = generate_mappings(tuples) bulk_insert = SalesforceBulkOperator( salesforce_conn_id='salesforce_conn_id', task_id='bulk_insert', operation='insert', object_name='Attachment', payload=mappings, external_id_field='Id', batch_size=10000, use_serial=False, ) mappings >> bulk_insert
备注:内容来源于stack exchange,提问作者displayname
相关产品推荐
相关产品推荐

