Airflow集成字节返回函数时报错:'bytes'对象无'__module__'属性
问题原因
Airflow的PythonOperator会自动将任务函数的返回值推送到XCom(跨任务通信系统),而XCom要求存储的数据必须是可序列化的类型(如字符串、字典、列表等)。你的函数返回的response.content是bytes类型,Airflow无法对其进行序列化,因此抛出'bytes' object has no attribute '__module__'错误。
解决方案
根据你的需求选择以下方案:
方案1:不需要保留返回值(直接处理字节内容)
如果后续任务不需要用到这个ZIP文件的字节内容,直接在函数内处理(比如保存到本地文件、对象存储),不返回bytes或返回None:
def run_scrape_latest_drug(): response = requests.get(zip_file_url) if response.status_code == 200: # 示例:将ZIP内容保存到本地 with open("/path/to/your/drug_data.zip", "wb") as f: f.write(response.content) return None # 显式返回None,避免自动返回bytes data = scrape_latest_drug() return data
方案2:需要将字节内容传递给后续任务
如果后续任务要使用这个ZIP内容,将bytes转换为可序列化的Base64字符串,后续任务再解码还原:
步骤1:修改抓取函数返回Base64字符串
import base64 def run_scrape_latest_drug(): response = requests.get(zip_file_url) if response.status_code == 200: # 将bytes转为Base64编码的字符串 return base64.b64encode(response.content).decode("utf-8") data = scrape_latest_drug() return data
步骤2:后续任务解码还原字节内容
import base64 from airflow.operators.python import PythonOperator def process_zip_file(ti): # 从XCom拉取Base64字符串 zip_base64_str = ti.xcom_pull(task_ids="scrape") # 解码回bytes zip_bytes = base64.b64decode(zip_base64_str) # 后续处理逻辑,比如解压ZIP等 pass # 在DAG中定义处理任务 task_process = PythonOperator( task_id="process_zip", python_callable=process_zip_file ) # 设置任务依赖 task_scrape >> task_process
注意事项
- 尽量避免通过XCom传递大文件内容,大文件推荐直接存储到外部存储(如S3、HDFS、本地文件系统),XCom仅传递文件路径或标识即可,避免占用Airflow元数据库资源。
内容的提问来源于stack exchange,提问作者ford_lnwza
相关产品推荐
相关产品推荐

