在Airflow中编写DAG提取余额总和时遇pipeline必须为列表错误
解决Apache Airflow连接MongoDB聚合时的TypeError问题
问题概述
在Airflow中编写DAG提取MongoDB钱包余额总和时,反复触发TypeError: pipeline must be a list错误,导致任务失败。
错误日志
TypeError: pipeline must be a list [2022-11-28, 11:24:51 UTC] {taskinstance.py:1401} INFO - Marking task as FAILED. dag_id=WALLET_DAGS, task_id=test_wallet, execution_date=20221128T112443, start_date=20221128T112449, end_date=20221128T112451
用户原实现代码
import logging import json logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) from datetime import datetime from airflow import DAG from airflow.operators.python import PythonOperator from airflow import AirflowException # Connection from airflow.providers.mongo.hooks.mongo import MongoHook def wallet_bal(): hook = MongoHook(mongo_conn_id='mongo_default') client = hook.get_conn() print(f"** Connected to MongoDB **- {client.server_info()}") db = client.test query = db.wallets.aggregate([{'$group': {'_id': None, 'count': {'$sum': '$balance.value'}}}]) wallets_collection = db.aggregate("wallets", query=query) logger.info(wallets_collection) logger.info("HI SEE THE PIPELINE") return 'MONGO DATA TO CONNECT' dag = DAG('WALLET_DAGS', description='Mongo Viewer',schedule_interval=None,start_date=datetime(2017, 3, 20), catchup=False) connect_case_mongo = PythonOperator(task_id='test_wallet', python_callable=wallet_bal, dag=dag) connect_case_mongo
问题分析与修复方案
错误核心是代码对MongoDB聚合操作的调用逻辑错误:
- 原代码先调用
db.wallets.aggregate()得到游标对象,再将游标传给db.aggregate(),但db.aggregate()要求传入的是聚合管道列表,而非游标。 db.aggregate()的使用不符合MongoDB API规范,正确做法是直接对目标集合调用聚合方法,或使用MongoHook封装的接口。
修复方式一:直接调用集合的aggregate方法
import logging import json logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) from datetime import datetime from airflow import DAG from airflow.operators.python import PythonOperator from airflow import AirflowException from airflow.providers.mongo.hooks.mongo import MongoHook def wallet_bal(): hook = MongoHook(mongo_conn_id='mongo_default') client = hook.get_conn() print(f"** Connected to MongoDB **- {client.server_info()}") db = client.test # 定义聚合管道列表 pipeline = [{'$group': {'_id': None, 'total_balance': {'$sum': '$balance.value'}}}] # 直接对wallets集合执行聚合,将游标转为列表读取结果 result = list(db.wallets.aggregate(pipeline)) if result: total_balance = result[0]['total_balance'] logger.info(f"总钱包余额: {total_balance}") else: logger.info("未找到钱包数据") return 'MONGO DATA TO CONNECT' dag = DAG('WALLET_DAGS', description='Mongo Viewer', schedule_interval=None, start_date=datetime(2017, 3, 20), catchup=False) connect_case_mongo = PythonOperator(task_id='test_wallet', python_callable=wallet_bal, dag=dag) connect_case_mongo
修复方式二:使用MongoHook封装的aggregate方法
import logging import json logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) from datetime import datetime from airflow import DAG from airflow.operators.python import PythonOperator from airflow import AirflowException from airflow.providers.mongo.hooks.mongo import MongoHook def wallet_bal(): hook = MongoHook(mongo_conn_id='mongo_default') # 定义聚合管道列表 pipeline = [{'$group': {'_id': None, 'total_balance': {'$sum': '$balance.value'}}}] # 通过MongoHook直接执行聚合,指定数据库和集合 result = list(hook.aggregate(database='test', collection='wallets', pipeline=pipeline)) if result: total_balance = result[0]['total_balance'] logger.info(f"总钱包余额: {total_balance}") else: logger.info("未找到钱包数据") return 'MONGO DATA TO CONNECT' dag = DAG('WALLET_DAGS', description='Mongo Viewer', schedule_interval=None, start_date=datetime(2017, 3, 20), catchup=False) connect_case_mongo = PythonOperator(task_id='test_wallet', python_callable=wallet_bal, dag=dag) connect_case_mongo
关键修复点
- 确保聚合操作传入的是管道列表(如
[{'$group': ...}]),而非游标或其他类型。 - 避免重复调用聚合方法,直接对目标集合执行聚合即可。
- 使用
list()将聚合游标转换为可读取的结果列表,方便提取余额总和。
内容的提问来源于stack exchange,提问作者Shiva
相关产品推荐
相关产品推荐

