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

Airflow(Google Cloud Composer):bytes类型无法JSON序列化报错

解决MySqlToGoogleCloudStorageOperator抛出的"Object of type 'bytes' is not JSON serializable"错误

问题分析

从你的代码和错误栈来看,你在尝试用Airflow的MySqlToGoogleCloudStorageOperator把dag_run表的数据导出到GCS时,碰到了JSON序列化失败的问题。根源很明确:查询结果里有bytes类型的字段,而标准JSON序列化器没法处理bytes对象。

你的任务代码

def gen_export_table_task(table_config):
    export_task = MySqlToGoogleCloudStorageOperator(
        task_id='export_dag_run_to_gcs',
        dag=dag,
        sql='SELECT * FROM dag_run LIMIT 10',
        bucket='gs://aca-composer-bq-test-vishesh',
        filename="cloudsql_to_bigquery_file.json",
        google_cloud_storage_conn_id='google_cloud_storage_default',
        mysql_conn_id='airflow_db'
    )
    return export_task

完整错误日志

[2019-12-19 11:32:29,649] {models.py:1796} ERROR - Object of type 'bytes' is not JSON serializable
Traceback (most recent call last):
File "/usr/local/lib/airflow/airflow/models.py", line 1659, in _run_raw_task
result = task_copy.execute(context=context)
File "/usr/local/lib/airflow/airflow/contrib/operators/mysql_to_gcs.py", line 106, in execute
files_to_upload = self._write_local_data_files(cursor)
File "/usr/local/lib/airflow/airflow/contrib/operators/mysql_to_gcs.py", line 153, in _write_local_data_files
s = json.dumps(row_dict)
File "/opt/python3.6/lib/python3.6/json/init.py", line 231, in dumps
return _default_encoder.encode(obj)
File "/opt/python3.6/lib/python3.6/json/encoder.py", line 199, in encode
chunks = self.iterencode(o, _one_shot=True)
File "/opt/python3.6/lib/python3.6/json/encoder.py", line 257, in iterencode
return _iterencode(o, 0)
File "/opt/python3.6/lib/python3.6/json/encoder.py", line 180, in default
o.class.name
TypeError: Object of type 'bytes' is not JSON serializable

为什么会出现这个问题?

Airflow的这个Operator在处理查询结果时,会把每一行转成字典,然后直接用json.dumps序列化。而dag_run表中的conf字段是BLOB类型(用来存储JSON配置的二进制形式),MySQL驱动会把它返回成bytes对象,这就触发了序列化错误。

两种可行的解决方法

方法一:在SQL层面转换二进制字段(最简单推荐)

直接修改你的SELECT语句,把二进制字段转换成JSON能识别的字符串格式。MySQL提供了TO_BASE64()函数可以把二进制数据编码成Base64字符串,这样导出后还能反向解码还原:

SELECT 
    dag_id,
    run_id,
    execution_date,
    start_date,
    end_date,
    state,
    TO_BASE64(conf) AS conf,  -- 把BLOB类型的conf转成Base64字符串
    run_type,
    external_trigger
FROM dag_run LIMIT 10

如果是其他字符串字段因为编码问题返回bytes,可以用CONVERT()指定编码:

SELECT CONVERT(your_column USING utf8mb4) AS your_column FROM dag_run LIMIT 10

方法二:自定义Operator处理序列化逻辑

如果不想修改SQL,可以继承原Operator,重写它的_write_local_data_files方法,添加对bytes类型的处理逻辑:

import json
from airflow.contrib.operators.mysql_to_gcs import MySqlToGoogleCloudStorageOperator

class BytesFriendlyMySqlToGCSOperator(MySqlToGoogleCloudStorageOperator):
    def _write_local_data_files(self, cursor):
        row_counter = 0
        files_to_upload = []
        tmp_file = self._local_tmp_handle()

        # 定义一个处理bytes的函数
        def serialize_value(value):
            if isinstance(value, bytes):
                # 可选:要么解码成字符串(如果是文本二进制),要么转Base64
                try:
                    return value.decode('utf-8')
                except UnicodeDecodeError:
                    return value.hex()  # 解码失败就用十六进制表示
            return value

        for row in cursor:
            row_dict = dict(zip(cursor.keys(), row))
            # 遍历字典,处理所有bytes类型的值
            processed_row = {k: serialize_value(v) for k, v in row_dict.items()}
            s = json.dumps(processed_row)
            tmp_file.write(f"{s}\n".encode('utf-8'))
            row_counter += 1
            if row_counter >= self.max_file_size:
                files_to_upload.append(tmp_file)
                tmp_file = self._local_tmp_handle()
                row_counter = 0
        if row_counter > 0:
            files_to_upload.append(tmp_file)
        return files_to_upload

# 替换原Operator为自定义的版本
def gen_export_table_task(table_config):
    export_task = BytesFriendlyMySqlToGCSOperator(
        task_id='export_dag_run_to_gcs',
        dag=dag,
        sql='SELECT * FROM dag_run LIMIT 10',
        bucket='gs://aca-composer-bq-test-vishesh',
        filename="cloudsql_to_bigquery_file.json",
        google_cloud_storage_conn_id='google_cloud_storage_default',
        mysql_conn_id='airflow_db'
    )
    return export_task

验证步骤

  1. 先查一下dag_run表的结构,确认二进制字段:
DESCRIBE dag_run;
  1. 先在MySQL客户端执行修改后的SQL,确认返回的都是字符串类型,再放到Airflow任务中测试。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 09:17:35