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

PyFlink结合AWS Kinesis SQL连接器使用UDTF时遇CAST函数不支持异常

在使用PyFlink Table API读写AWS Kinesis流时,不使用UDTF功能一切正常,但添加带@udtf装饰器的扁平化解构函数后,执行INSERT INTO SQL语句时触发如下异常:

pyflink.util.exceptions.TableException: org.apache.flink.table.api.TableException: Unsupported Python SqlFunction CAST.

核心代码片段:

@udtf(result_types=[DataTypes.STRING(), DataTypes.INT()])
def flatten_row(row: Row) -> Row:
   for s in row["members"]:
       yield Row(str(s["id"]), s["name"])
result_table = input_table.flat_map(flatten_row).alias("id", "name")
table_env.create_temporary_view("result_table", result_table)
table_result = table_env.execute_sql(f"INSERT INTO {output_table_name} SELECT * FROM result_table")

解决方案

  • 消除隐式类型转换
    检查输出Kinesis表的Schema,确保UDTF的result_types与表字段类型完全一致,避免Flink触发自动CAST操作。同时移除UDTF函数内的手动类型转换,直接返回与目标类型匹配的原始数据:

    @udtf(result_types=[DataTypes.STRING(), DataTypes.INT()])
    def flatten_row(row: Row):
        for s in row["members"]:
            # 确保s["id"]为STRING类型,s["name"]为INT类型,与result_types和输出表Schema严格匹配
            yield s["id"], s["name"]
    
  • 改用Table API原生insert_into方法
    避免通过创建临时视图再执行SQL INSERT的方式,直接使用Table API的insert_into方法提交任务,绕开SQL层的类型转换逻辑:

    result_table = input_table.flat_map(flatten_row).alias("id", "name")
    # 直接执行插入操作,无需通过execute_sql
    result_table.insert_into(output_table_name).execute()
    
  • 升级Flink版本(可选)
    Flink 1.16版本在Python UDTF与SQL交互的类型处理上存在已知兼容性问题,升级到1.17及以上版本可修复部分相关BUG。


注意事项

  • 确认输入表的members字段类型定义正确(例如ARRAY<ROW<id STRING, name INT>>)
  • 确保输出Kinesis表的字段名称、类型、nullable属性与UDTF输出完全匹配

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 18:05:17