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

如何将本地递归Python函数改写为可处理BQ数据的Airflow任务

Airflow 函数迁移适配方案

问题诊断

你当前的Airflow版本代码存在核心逻辑错误,是把两层函数的代码混淆堆叠导致的:

  1. 本地版本中map_manufacturer_model是逐行处理单个序列号的小函数,仅用于pandas的apply调用
  2. 你在写Airflow版本时,把BQ数据拉取、全局处理逻辑都塞进了这个逐行处理的函数里,还保留了入参s,同时函数外还写了apply调用代码,逻辑完全混乱
  3. 你处理完的DataFrame没有回写BigQuery,任务跑完数据就会丢失,等于没有执行有效处理
  4. 补充说明:你的代码没有递归逻辑,Airflow执行普通Python函数的逻辑和本地完全一致,不需要做特殊的递归适配。

修正后可运行代码

import pandas as pd
from airflow import models
from airflow.operators.python import PythonOperator
from airflow.providers.google.cloud.hooks.bigquery import BigQueryHook

# 全局映射常量,放在函数外避免重复初始化
manufacturers = {'G4F0': 'FLN', 'G4F1': 'FLN', 'G4F9': 'FLN', 'G4K0': 'HWL', 'E6S1': 'LPG', 'E6S2': 'LPG'}
meter_models = {'G4F0': {'1': 'G4SZV-1', '2': 'G4SZV-2'},
                'G4F9': {'': 'G4SZV-1'},
                'G4F1': {'': 'G4SDZV-2'},
                'G4K0': {'': 'BK-G4E'},
                'E6S1': {'': 'E6VG470'},
                'E6S2': {'': 'E6VG470'},
                }

def process_meter_data():
    # 1. 用Airflow官方BQ Hook拉取数据,不需要自己维护客户端密钥,权限统一走Airflow连接配置
    bq_hook = BigQueryHook(gcp_conn_id='your_gcp_connection_id', use_bqstorage_api=True)
    query_string = """
    SELECT * FROM `your_project.your_dataset.bq_table`
    """
    gas_data = bq_hook.get_pandas_df(query_string)

    # 2. 定义逐行处理的内部函数,仅负责单个序列号的映射
    def map_manufacturer_model(s):
        s = str(s)
        model = ''
        try:
            manufacturer = manufacturers[s[:4]]
            for k, m in meter_models[s[:4]].items():
                if s[-4:].startswith(k):
                    model = m
                    break
        except KeyError:
            manufacturer = ''
        return pd.Series({'New Meter Manufacturer': manufacturer,
                          'New Meter Model': model
                          })
    
    # 3. 批量处理数据
    gas_data[['New Meter Manufacturer', 'New Meter Model']] = gas_data['New Serial Number'].apply(map_manufacturer_model)

    # 4. 处理完的数据回写BQ目标表
    bq_hook.insert_rows_from_dataframe(
        table='your_project.your_dataset.target_bq_table',
        dataframe=gas_data,
        replace=True  # 根据需求选择追加还是覆盖
    )

# DAG定义
with models.DAG(
    'test_dag',
    schedule_interval='0 8 * * *',
    default_args=default_dag_args, # 这里替换为你自己定义的default_args即可
    catchup=False
) as dag:
    map_manufacturer_model_task = PythonOperator(
        task_id='map_manufacturer_model_function',
        python_callable=process_meter_data
    )

最佳实践建议

  • 如果数据量超过10万行,不建议拉到Airflow Worker本地用pandas处理,直接把映射逻辑写成BigQuery SQL UDF,在BQ端直接完成转换,性能是Python处理的几十倍,还不会出现Worker内存不足的问题
  • 映射常量如果后续会频繁变更,可以存在BQ小表或Airflow Variable里,不要硬编码在DAG代码中,方便后续调整
  • 处理完的数据如果需要给下游任务使用,直接读取回写后的BQ表即可,不需要在Airflow任务之间传递DataFrame

内容的提问来源于stack exchange,提问作者Mizanur Choudhury

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.23 20:54:01