lydtcsdn 2020-05-21 16:54 采纳率: 0%
浏览 321

求教flink自定义python udf时TIMESTAMP类型问题?

版本:pyflink1.10
pyflink使用python udf的时候数据类型定义为DataTypes.TIMESTAMP(),但是执行时被解释为long型,报not match
源码:

@udf(input_types=[DataTypes.INT(), DataTypes.INT(), DataTypes.TIMESTAMP()], result_type=DataTypes.INT())
def add_new(i, j, k):
    return i + j    #k没用我就是测试一下

报错:

py4j.protocol.Py4JJavaError: An error occurred while calling o80.select.
: org.apache.flink.table.api.ValidationException: Given parameters do not match any signature. 
Actual: (java.lang.Integer, java.lang.Integer, java.time.LocalDateTime) 
Expected: (int, int, long)

补充:使用kafka数据源,kafka版本为0.11.3
source部分:

st_env.from_path("source") \
.select("a, add_new(b, c, rowtime) as sumbc") \
.insert_into("sink")

schema解析:

.with_schema(  
    Schema()
     .field("a", DataTypes.STRING())
     .field("b", DataTypes.INT())
     .field("c", DataTypes.INT())
     .field("time", DataTypes.TIMESTAMP())
     .rowtime(
                    Rowtime()
                     .timestamps_from_field("time")
                     .watermarks_periodic_bounded(60000)
                      )
     ) 

kafka单条数据样式:

{'a': 'a', 'b': 1, 'c': 1, 'time': '2013-01-01T00:14:13Z'}

或者是哪里写错了吗?请各位指教

另外当使用register_table_source创建表时DataTypes.TIMESTAMP()的类型会显示为java.sql.Timestamp,这种类型不一致会造成问题吗?

  • 写回答

1条回答 默认 最新

报告相同问题?

悬赏问题

  • ¥15 安卓adb backup备份应用数据失败
  • ¥15 eclipse运行项目时遇到的问题
  • ¥15 关于#c##的问题:最近需要用CAT工具Trados进行一些开发
  • ¥15 南大pa1 小游戏没有界面,并且报了如下错误,尝试过换显卡驱动,但是好像不行
  • ¥15 没有证书,nginx怎么反向代理到只能接受https的公网网站
  • ¥50 成都蓉城足球俱乐部小程序抢票
  • ¥15 yolov7训练自己的数据集
  • ¥15 esp8266与51单片机连接问题(标签-单片机|关键词-串口)(相关搜索:51单片机|单片机|测试代码)
  • ¥15 电力市场出清matlab yalmip kkt 双层优化问题
  • ¥30 ros小车路径规划实现不了,如何解决?(操作系统-ubuntu)