升级GCP Composer与Airflow后,MSSQLToGCSOperator不生成空文件问题
解决MSSQLToGCSOperator无数据时不生成空文件的问题
升级GCP Composer(1.19.13)和Airflow(2.3.3)后,遇到一个问题:之前版本里,当MSSQL查询无返回数据时,MSSQLToGCSOperator会在GCS存储桶里创建空文件,但现在不会了,导致后续依赖这个文件的任务报“文件未找到”错误。
我当前使用的代码如下:
from airflow.providers.google.cloud.transfers.mssql_to_gcs import MSSQLToGCSOperator mssql_to_gcs = MSSQLToGCSOperator( task_id='MYSQL_TO_GCS_{0}'.format(TABLE_NAME), mssql_conn_id='con_mssql_dba_prd', gcp_conn_id='google_cloud_storage_default', sql='select_{0}.sql'.format(TABLE_NAME), bucket=SOURCE_BUCKET, filename='composer/{0}/gcp_{0}{1}.json'.format(TABLE_NAME, DATE_FORMAT), dag=dag )
三种可行解决方案
1. 自定义Operator强制生成空文件
继承原Operator重写逻辑,查询完后检查文件是否存在,不存在就创建空文件:
from airflow.providers.google.cloud.transfers.mssql_to_gcs import MSSQLToGCSOperator from google.cloud import storage from airflow.providers.google.cloud.hooks.google_cloud import GoogleCloudHook class CustomMSSQLToGCSOperator(MSSQLToGCSOperator): def execute(self, context): # 执行原Operator的导出逻辑 super().execute(context) # 初始化GCS客户端 hook = GoogleCloudHook(gcp_conn_id=self.gcp_conn_id) storage_client = storage.Client(credentials=hook.get_credentials()) bucket = storage_client.bucket(self.bucket) blob = bucket.blob(self.filename) # 检查文件是否存在,不存在则上传空内容 if not blob.exists(): blob.upload_from_string('')
之后直接用这个自定义Operator替换原有的即可,参数不用改。
2. 后续任务前置检查并创建空文件
如果不想改Operator,可以加一个PythonOperator在后续任务之前,负责检查文件,不存在就创建:
from airflow.operators.python import PythonOperator from google.cloud import storage from airflow.providers.google.cloud.hooks.google_cloud import GoogleCloudHook def ensure_file_exists(bucket_name, file_path, gcp_conn): hook = GoogleCloudHook(gcp_conn_id=gcp_conn) client = storage.Client(credentials=hook.get_credentials()) bucket = client.bucket(bucket_name) blob = bucket.blob(file_path) if not blob.exists(): blob.upload_from_string('') # 创建检查任务 ensure_file_task = PythonOperator( task_id='ensure_gcs_file_exists', python_callable=ensure_file_exists, op_kwargs={ 'bucket_name': SOURCE_BUCKET, 'file_path': 'composer/{0}/gcp_{0}{1}.json'.format(TABLE_NAME, DATE_FORMAT), 'gcp_conn': 'google_cloud_storage_default' }, dag=dag ) # 调整任务依赖:导出任务 → 检查创建任务 → 后续任务 mssql_to_gcs >> ensure_file_task >> your_downstream_task
3. 修改SQL查询返回空行
直接改你的select_{0}.sql查询语句,确保即使原查询无数据,也返回一行空值:
-- 原查询 SELECT col1, col2, col3 FROM your_target_table -- 追加UNION ALL,当原查询无结果时返回空行 UNION ALL SELECT NULL, NULL, NULL WHERE NOT EXISTS (SELECT 1 FROM your_target_table)
注意要把NULL的数量和顺序对应上原查询的字段数,这样Operator会生成包含这行空数据的JSON文件,后续任务就能找到文件了。
内容的提问来源于stack exchange,提问作者Benito De la Cruz
相关产品推荐
相关产品推荐

