如何将本地递归Python函数改写为可处理BQ数据的Airflow任务
Airflow 函数迁移适配方案
问题诊断
你当前的Airflow版本代码存在核心逻辑错误,是把两层函数的代码混淆堆叠导致的:
- 本地版本中
map_manufacturer_model是逐行处理单个序列号的小函数,仅用于pandas的apply调用 - 你在写Airflow版本时,把BQ数据拉取、全局处理逻辑都塞进了这个逐行处理的函数里,还保留了入参
s,同时函数外还写了apply调用代码,逻辑完全混乱 - 你处理完的DataFrame没有回写BigQuery,任务跑完数据就会丢失,等于没有执行有效处理
- 补充说明:你的代码没有递归逻辑,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
相关产品推荐
相关产品推荐

