You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

升级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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.09 23:50:29