【发布时间】:2021-12-09 19:57:24
【问题描述】:
我有一个 flink 表,比如CREATE TABLE source(id int, name string) with (...),还有一个目标表,比如CREATE TABLE destination(id int, unique_name string) with (...)。 unique_name是使用内部flink流程函数中的业务逻辑计算出来的。
所以我们可以安全地假设源模式将与目标模式完全相同(名称和数据类型)。
我使用数据流 API 做了一些低级别的 process 以获取 destination 数据流。它有outputType 和GenericType<org.apache.flink.types.Row>。当我再次将destination 数据流转换回表时,出现以下错误。
org.apache.flink.table.api.ValidationException: Column types of query result and sink
for registered table 'default_catalog.default_database.destination' do not match.
Cause: Different number of columns.
Query schema: [f0: RAW('org.apache.flink.types.Row', '...')]
Sink schema: [id: INT, name: STRING]
虽然我可以使用以下代码解决此问题,但我想将其泛化并从目标 Table 获取 RowTypeInformation。有没有办法从flinkTable获取TypeInformation。
tableEnv.fromDataStream(destionationDataStream.map(x -> x).returns(Types.ROW(Types.Int, Types.String))
【问题讨论】:
标签: apache-flink flink-streaming