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

