Airflow(Google Cloud Composer):bytes类型无法JSON序列化报错
问题分析
从你的代码和错误栈来看,你在尝试用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
验证步骤
- 先查一下
dag_run表的结构,确认二进制字段:
DESCRIBE dag_run;
- 先在MySQL客户端执行修改后的SQL,确认返回的都是字符串类型,再放到Airflow任务中测试。
内容的提问来源于stack exchange,提问作者DataVishesh

