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

Airflow Python Callable调用函数报JSON不可序列化错误如何解决

问题根因

你遇到的两个报错直接关联,PyCharm的类型提示不是误报:

Expected type 'str', got '() -> Union[str, Dict[str, Union[str, Any]]]' instead

这个提示明确说明你在某个要求字符串/可序列化对象的参数位置,传入了函数对象本身,而非函数执行后的返回值。Airflow在解析DAG、序列化任务参数或XCom返回值时,会尝试对所有传入的参数做JSON序列化,函数对象无法被JSON序列化,就触发了你看到的TypeError: Object of type function is not JSON serializable错误。

修复步骤
  • 检查PythonOperator的参数配置
    优先核对两处写法:
    1. 如果你是把这个自定义函数作为python_callable参数传入,注意只传函数对象本身,不要加括号调用:
    # 正确写法
    def query_cluster_status():
        # 你的业务逻辑
        if cluster_available:
            return {"Cluster": cluster_name, "Status": cluster_status, "Restored": "Yes"}
        else:
            return {"Status": cluster_status, "Restored": "No"}
    
    check_cluster_task = PythonOperator(
        task_id="check_cluster",
        python_callable=query_cluster_status,  # 不要加括号写成query_cluster_status()
        do_xcom_push=True
    )
    
    1. 如果你是把这个自定义函数的结果作为参数传给其他算子,必须调用函数(加括号)获取返回值,不要直接传函数对象:
    # 错误写法,会触发序列化报错
    downstream_task = PythonOperator(
        task_id="process_cluster",
        python_callable=process_func,
        op_kwargs={"cluster_info": query_cluster_status}  # 传了函数对象,未调用
    )
    
    # 正确写法
    downstream_task = PythonOperator(
        task_id="process_cluster",
        python_callable=process_func,
        op_kwargs={"cluster_info": query_cluster_status()}
    )
    
  • 检查返回值是否包含不可序列化元素
    确认你的自定义函数返回的字典/字符串中,没有嵌套函数、类实例、AWS连接对象这类不可序列化的元素,你可以在本地单独执行函数,用json.dumps()方法测试返回值是否可被序列化:
    import json
    res = query_cluster_status()
    print(json.dumps(res)) # 不报错说明返回值本身没有序列化问题
    
  • 适配低版本Airflow的序列化规则
    如果你使用的是Airflow 1.x版本,XCom对复杂字典的序列化支持不完善,可以手动把返回值转成JSON字符串再返回,下游取值时再做JSON解析即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 11:09:01