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

Airflow v1.8.0中MySQL转GCS出现byte-like object required错误求助

问题分析与解决方案

这个错误是Python 3 环境下 Airflow v1.8.0 中 MySQLToGoogleCloudStorageOperator 的兼容性问题,具体原因如下:

在 Python 3 里,tempfile.NamedTemporaryFile 默认是以**二进制模式('wb')**创建文件句柄的,但 Airflow v1.8.0 的这个 Operator 在_write_local_data_files方法中,直接用json.dump往二进制文件句柄写入字符串类型的JSON数据——而二进制文件只接受字节对象,这就触发了TypeError: a bytes-like object is required, not 'str'异常。

从报错堆栈也能明确看到,错误卡在json.dump(row_dict, tmp_file_handle)这一步,底层调用fp.write(chunk)时,因为文件是二进制模式,无法写入字符串才抛出了错误。


解决方法

1. 升级Airflow版本(推荐)

Airflow 在 v1.9.0 及之后的版本已经修复了这个 Python 3 兼容性问题,直接升级到新版本是最彻底的解决方案,同时还能获得更多功能和 bug 修复:

pip install --upgrade apache-airflow==<你的目标版本,比如1.10.15>

(如果需要保持版本兼容性,也可以选择1.9.x系列的稳定版本)

2. 临时修改Airflow源码(快速应急)

如果暂时无法升级,可以直接修改MySQLToGoogleCloudStorageOperator的源码:
找到你的Airflow安装路径下的airflow/contrib/operators/mysql_to_gcs.py文件,定位到_write_local_data_files方法中创建临时文件的代码行:

# 原代码
tmp_file = tempfile.NamedTemporaryFile(delete=False)

修改为指定文本模式和编码:

# 修改后
tmp_file = tempfile.NamedTemporaryFile(mode='w', delete=False, encoding='utf-8')

修改后重启Airflow调度器即可生效。

3. 自定义Operator(无权限修改源码时)

如果没有权限修改系统级的Airflow源码,可以自己实现一个修正版的Operator:

import tempfile
import json
from airflow.contrib.operators.mysql_to_gcs import MySQLToGoogleCloudStorageOperator

class FixedMySQLToGCSOperator(MySQLToGoogleCloudStorageOperator):
    def _write_local_data_files(self, cursor):
        """修正了Python3下的文件模式问题"""
        schema = []
        for field in cursor.description:
            schema.append({
                'name': field[0],
                'type': self.type_map[field[1]],
                'mode': 'NULLABLE',
            })

        files_to_upload = []
        tmp_file_handle = tempfile.NamedTemporaryFile(mode='w', delete=False, encoding='utf-8')
        for row in cursor:
            row_dict = dict(zip((field[0] for field in cursor.description), row))
            json.dump(row_dict, tmp_file_handle)
            tmp_file_handle.write('\n')
        tmp_file_handle.close()
        files_to_upload.append(tmp_file_handle.name)

        if self.schema_filename:
            schema_file_handle = tempfile.NamedTemporaryFile(mode='w', delete=False, encoding='utf-8')
            json.dump({'fields': schema}, schema_file_handle)
            schema_file_handle.close()
            files_to_upload.append(schema_file_handle.name)

        return files_to_upload

然后在你的DAG中替换使用这个自定义Operator:

export_waybills = FixedMySQLToGCSOperator( 
    task_id='extract_waybills', 
    mysql_conn_id = 'podiotestmySQL', 
    sql = 'SELECT * FROM podiodb.logistics_waybills', 
    bucket='podio-reader-storage', 
    filename= 'podio-data/waybills{}.json', 
    schema_filename='podio-data/schema/waybills.json', 
    dag=dag
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 07:15:12