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

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

问题

在Airflow中编写了名为data-sync的DAG,其中get_view_buffer任务负责处理CSV数据并返回StringIO对象。任务执行过程本身正常,但最终抛出AttributeError,提示'_io.StringIO'对象没有'__module__'属性。

DAG代码

from datetime import timedelta
import pendulum
from airflow.decorators import dag
from stage import stage_data
from table_async_pg import run_async
from read_mem_view import get_view_buffer
@dag(
dag_id = "data-sync",
schedule_interval = '*/30 * * * *',
start_date=pendulum.datetime(2023, 3, 9, tz="Asia/Hong_Kong"),
catchup=False,
dagrun_timeout=timedelta(minutes=20),
)
def Pipeline():
    CSV_URL="myurl"
    a = stage_data()
    b = run_async()
    c = get_view_buffer(CSV_URL)

    [a,c] >> b


pipeline = Pipeline()

任务代码

@task
def get_view_buffer(URL):
    SKIPROWS=2
    p = Path(r'/home/d5291029/cert/truststore-prod-2.0.7.pem')
    if p.exists():
        with requests.Session() as s:
            start = time.time()
            download = s.get(URL,verify=p)
            decoded_content = download.content.decode('utf-8')
            end = time.time() - start
            print(f'time it takes to download in seconds: {end}')
            reader = csv.reader(decoded_content.splitlines(),delimiter=',')
            output = io.StringIO()
            writer = csv.writer(output,lineterminator='\n')
            for row in reader:
                if reader.line_num > SKIPROWS:
                    writer.writerow(row)
            
            output.seek(0)
            return output

报错信息

[2023-03-10, 17:30:07 HKT] {python.py:177} INFO - Done. Returned value was: <_io.StringIO object at 0x7f2a70fdfee0>
[2023-03-10, 17:30:07 HKT] {taskinstance.py:1768} ERROR - Task failed with exception
..................
  File "/opt/.venv/lib/python3.9/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需要将任务返回值序列化后存储到元数据库,但StringIO对象无法被Airflow的JSON序列化器处理。

修改get_view_buffer任务的返回逻辑,不直接返回StringIO对象,而是返回其包含的字符串内容:

# 替换原有的return output语句
return output.getvalue()

如果后续任务需要使用StringIO对象,可以在接收返回值的任务中,将字符串重新转换为StringIO:

@task
def run_async(csv_content):
    output = io.StringIO(csv_content)
    # 后续业务处理逻辑

内容的提问来源于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 19:33:27