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

Airflow任务返回io.StringIO对象触发AttributeError错误求助

解决Airflow 2.5.1中返回io.StringIO对象导致的XCom序列化错误

环境信息

使用Airflow官方Docker镜像2.5.1,核心配置如下:

  • 执行器:CeleryExecutor
  • 元数据库:PostgreSQL
  • Python版本:3.7.16

问题代码

编写的DAG定义:

import datetime
import pendulum
import io
from airflow.decorators import dag, task


@dag(
    dag_id="my_beauty",
    schedule_interval="0 0 * * *",
    start_date=pendulum.datetime(2023, 3, 13, tz="UTC"),
    catchup=False,
    dagrun_timeout=datetime.timedelta(minutes=60),
)
def MyBeauty():
    @task
    def get_buffer():
        content = io.StringIO("some words here blabla")
        return content
    
    a = get_buffer()

dag = MyBeauty()

报错信息

运行任务时触发如下错误:

Traceback (most recent call last):
  File "/home/airflow/.local/bin/airflow", line 8, in <module>
    sys.exit(main())
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/__main__.py", line 39, in main
    args.func(args)
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/cli/cli_parser.py", line 52, in command
    return func(*args, **kwargs)
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/utils/cli.py", line 108, in wrapper
    return f(*args, **kwargs)
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/cli/commands/task_command.py", line 575, in task_test
    ti.run(ignore_task_deps=True, ignore_ti_state=True, test_mode=True)
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/utils/session.py", line 75, in wrapper
    return func(*args, session=session, **kwargs)
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/models/taskinstance.py", line 1670, in run
    mark_success=mark_success, test_mode=test_mode, job_id=job_id, pool=pool, session=session
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/utils/session.py", line 72, in wrapper
    return func(*args, **kwargs)
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/models/taskinstance.py", line 1374, in _run_raw_task
    self._execute_task_with_callbacks(context, test_mode)
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/models/taskinstance.py", line 1520, in _execute_task_with_callbacks
    result = self._execute_task(context, task_orig)
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/models/taskinstance.py", line 1588, in _execute_task
    self.xcom_push(key=XCOM_RETURN_KEY, value=xcom_value, session=session)
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/utils/session.py", line 72, in wrapper
    return func(*args, **kwargs)
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/models/taskinstance.py", line 2297, in xcom_push
    session=session,
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/utils/session.py", line 72, in wrapper
    return func(*args, **kwargs)
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/models/xcom.py", line 240, in set
    map_index=map_index,
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/models/xcom.py", line 627, in serialize_value
    return json.dumps(value, cls=XComEncoder).encode("UTF-8")
  File "/usr/local/lib/python3.7/json/__init__.py", line 238, in dumps
    **kw).encode(obj)
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/utils/json.py", line 176, in encode
    return super().encode(o)
  File "/usr/local/lib/python3.7/json/encoder.py", line 199, in encode
    chunks = self.iterencode(o, _one_shot=True)
  File "/usr/local/lib/python3.7/json/encoder.py", line 257, in iterencode
    return _iterencode(o, 0)
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/utils/json.py", line 153, in default
    CLASSNAME: o.__module__ + "." + o.__class__.__qualname__,
AttributeError: '_io.StringIO' object has no attribute '__module__'

问题原因

Airflow的XCom默认仅支持JSON可序列化的数据类型。当任务返回_io.StringIO对象时,XCom的序列化编码器(XComEncoder)尝试读取对象的__module__属性完成序列化,但_io.StringIO对象没有该属性,导致报错。

解决方案

方案1:返回字符串内容(推荐)

直接返回StringIO对象中的字符串内容,而非对象本身,这是最简单且符合XCom设计的方式:

@task
def get_buffer():
    content = io.StringIO("some words here blabla")
    return content.getvalue()  # 返回字符串而非StringIO对象

方案2:自定义XCom序列化器(特殊场景)

如果必须在任务间传递流对象,可以自定义XCom的序列化/反序列化逻辑:

  1. 编写自定义编码器和解码器
from airflow.utils.json import XComEncoder
from io import StringIO

class CustomXComEncoder(XComEncoder):
    def default(self, o):
        if isinstance(o, StringIO):
            return {"__type__": "StringIO", "content": o.getvalue()}
        return super().default(o)

def custom_xcom_decoder(obj):
    if "__type__" in obj and obj["__type__"] == "StringIO":
        return StringIO(obj["content"])
    return obj
  1. 在airflow.cfg中配置自定义组件
[core]
xcom_encoder = your_dag_module.CustomXComEncoder
xcom_decoder = your_dag_module.custom_xcom_decoder

注意:需确保所有Airflow Worker都能加载到该自定义模块,适合有一定经验的用户。

方案3:避免跨任务传递(本地使用)

如果StringIO对象仅在当前任务内部使用,无需传递给其他任务,直接在任务内处理完毕即可,不要返回该对象,避免触发XCom序列化。

总结

优先选择方案1,因为XCom的定位是传递轻量、可序列化的数据,复杂对象并不适合通过XCom在任务间传递。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 02:30:01