CloudSqlExecuteQueryOperator无法向XCom返回结果的问题求助
我有一个Airflow流水线,用于向GCP Cloud SQL上的PostgreSQL数据库执行UPSERT操作。此前在本地数据库使用PostgresOperator时一切正常,但需将DAG执行迁移至搭配Cloud SQL的Cloud Composer中。我已完成复杂配置,确保Cloud Composer的Kubernetes集群Pod可连接数据库,且部分算子已成功创建数据库表。
但问题出在尝试从CloudSqlExecuteQueryOperator获取返回值时。我的模板化SQL命令如下:
WITH input_rows(plate_name, plate_type, number_wells) AS ( VALUES {% for row in ti.xcom_pull(task_ids='get_message_rows', key='return_value') %} ('{{ row['plate'] }}', 'culture', 384) {{ ",\n" if not loop.last else "" }} {% endfor %} ) , ins AS ( INSERT INTO plate (plate_name, plate_type, number_wells) SELECT * FROM input_rows ON CONFLICT (plate_name) DO NOTHING RETURNING plate_id, plate_name, plate_type, number_wells ) SELECT 'i' AS source , plate_id, plate_name, plate_type, number_wells FROM ins UNION ALL SELECT 's' AS source , c.plate_id, input_rows.plate_name, input_rows.plate_type, input_rows.number_wells FROM input_rows JOIN plate c USING (plate_name);
算子执行代码如下:
write_plates = CloudSQLExecuteQueryOperator( task_id="write_plates", sql=sql_plate_template, gcp_cloudsql_conn_id=connection_name, # environment variable AIRFLOW_CONN_<conn_id> ) # 该算子因获取到NoneType而失败 formatted_results = some_python_operator_with_decoration(write_plates.output) write_plates >> formatted_results
我能看到数据已成功插入表中,手动执行该查询也能得到预期结果。请问CloudSqlExecuteQueryOperator是否不支持向XCom返回结果?是否需要设置autocommit为true?有没有可行的解决办法?我期望能像使用PostgresOperator那样迭代处理返回的行和列,示例代码如下:
@task def some_python_operator_with_decoration(postgres_results): results = [] for row in postgres_results: results.append({'source': row[0], "id": row[1], "name": row[2]}) return results
CloudSqlExecuteQueryOperator默认不返回查询结果
该算子的设计逻辑是仅执行SQL语句,默认不会将查询结果推送到XCom。这和PostgresOperator的行为不同,PostgresOperator默认会返回查询结果并存储到XCom中。autocommit设置不影响结果返回
autocommit参数仅控制事务提交行为,和是否返回查询结果没有直接关联,即使设置为True也无法让CloudSqlExecuteQueryOperator返回结果。可行替代方案
有两种常用方法可以实现获取查询结果的需求:方法一:使用CloudSQLHook自定义Python函数
放弃CloudSqlExecuteQueryOperator,改用
CloudSQLHook在Python任务中手动执行查询并返回结果,示例代码如下:from airflow.providers.google.cloud.hooks.cloud_sql import CloudSQLHook from airflow.decorators import task @task def run_upsert_and_get_results(connection_name, sql_template, ti): # 渲染SQL模板 rendered_sql = ti.render_template(sql_template, {}) # 初始化CloudSQLHook hook = CloudSQLHook(gcp_cloudsql_conn_id=connection_name) # 执行查询并获取结果 with hook.get_conn() as conn: with conn.cursor() as cursor: cursor.execute(rendered_sql) results = cursor.fetchall() # 返回结果到XCom return results # 调用自定义任务 write_plates = run_upsert_and_get_results(connection_name, sql_plate_template) formatted_results = some_python_operator_with_decoration(write_plates.output) write_plates >> formatted_results方法二:修改CloudSqlExecuteQueryOperator的行为(不推荐)
如果一定要使用该算子,可以通过继承并重写
execute方法,添加结果推送逻辑。但这种方法需要维护自定义算子,后续Airflow版本更新可能会引发兼容性问题,示例代码如下:from airflow.providers.google.cloud.operators.cloud_sql import CloudSQLExecuteQueryOperator from airflow.utils.context import Context class CloudSQLExecuteQueryWithResultOperator(CloudSQLExecuteQueryOperator): def execute(self, context: Context): # 执行父类的execute逻辑 super().execute(context) # 获取连接并执行查询获取结果 hook = self._get_hook() with hook.get_conn() as conn: with conn.cursor() as cursor: cursor.execute(self.sql) results = cursor.fetchall() # 将结果推送到XCom ti = context['ti'] ti.xcom_push(key='return_value', value=results) return results # 使用自定义算子 write_plates = CloudSQLExecuteQueryWithResultOperator( task_id="write_plates", sql=sql_plate_template, gcp_cloudsql_conn_id=connection_name, ) formatted_results = some_python_operator_with_decoration(write_plates.output) write_plates >> formatted_results
内容的提问来源于stack exchange,提问作者CONNOR TRUEX

