【发布时间】:2020-10-27 17:02:51
【问题描述】:
我想实现一种动态 FTPSensor。使用贡献的 FTP 传感器,我设法使它以这种方式工作:
ftp_sensor = FTPSensor(
task_id="detect-file-on-ftp",
path="./data/test.txt",
ftp_conn_id="ftp_default",
poke_interval=5,
dag=dag,
)
它工作得很好。但我需要传递动态路径和 ftp_conn_id 参数。 IE。我在之前的任务和 ftp_sensor 任务中生成了一堆新连接,如果 FTP 上存在文件,我想检查我之前生成的每个新连接。
所以我想首先从 XCom 获取连接的 ID。 我从 XCom 中的上一个任务发送它们,但似乎我无法在任务之外访问 XCom。 例如。我的目标是:
active_ftp_connections = context['ti'].xcom_pull(key='active_ftps')
for conn in active_ftp_connections:
ftp_sensor = FTPSensor(
task_id="detect-file-on-ftp",
path=conn['path'],
ftp_conn_id=conn['connection'],
poke_interval=5,
dag=dag,
)
但这似乎不是一个可能的解决方案。
然后我浪费了大量时间尝试创建我的自定义 FTPSensor 以动态传递我需要的数据,但现在我得出的结论是我需要传感器和操作员之间的混合,因为我需要例如,保留戳功能,但也具有执行功能。 我想一种选择是编写一个自定义运算符,从传感器基类实现 poke,但现在可能太累了,无法尝试。
您知道如何实现我的目标吗?我似乎无法在互联网上找到有关该主题的任何材料 - 也许只有我一个人。 如果问题不清楚,请告诉我,以便我提供更多详细信息。
更新
我现在认为这是可能的
def get_active_ftps(**context):
active_ftp_connestions = context['ti'].xcom_pull(key='active_ftps')
return active_ftp_connestions
for ftp in get_active_ftps():
ftp_sensor = FTPSensor(
task_id="detect-file-on-ftp",
path="./"+ ftp['folder'] +"/test.txt",
ftp_conn_id=ftp['conn_id'],
poke_interval=5,
dag=dag,
)
但它会引发错误:Broken DAG: [/usr/local/airflow/dags/copy_file_from_ftp.py] 'ti'
【问题讨论】: