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条回答 默认 最新

报告相同问题?

悬赏问题

  • ¥60 求一个简单的网页(标签-安全|关键词-上传)
  • ¥35 lstm时间序列共享单车预测,loss值优化,参数优化算法
  • ¥15 基于卷积神经网络的声纹识别
  • ¥15 Python中的request,如何使用ssr节点,通过代理requests网页。本人在泰国,需要用大陆ip才能玩网页游戏,合法合规。
  • ¥100 为什么这个恒流源电路不能恒流?
  • ¥15 有偿求跨组件数据流路径图
  • ¥15 写一个方法checkPerson,入参实体类Person,出参布尔值
  • ¥15 我想咨询一下路面纹理三维点云数据处理的一些问题,上传的坐标文件里是怎么对无序点进行编号的,以及xy坐标在处理的时候是进行整体模型分片处理的吗
  • ¥15 CSAPPattacklab
  • ¥15 一直显示正在等待HID—ISP