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

在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聚合操作的调用逻辑错误:

  1. 原代码先调用db.wallets.aggregate()得到游标对象,再将游标传给db.aggregate(),但db.aggregate()要求传入的是聚合管道列表,而非游标。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 03:25:28