【问题标题】:Flink Table to DataStream: how to access column name?Flink Table to DataStream:如何访问列名?
【发布时间】:2021-03-30 13:47:29
【问题描述】:

我想使用 Flink SQL 将 Kafka 主题消费到表中,然后将其转换回 DataStream。

这里是SOURCE_DDL

CREATE TABLE kafka_source (
    user_id BIGINT,
    datetime TIMESTAMP(3),
    last_5_clicks STRING
) WITH (
    'connector' = 'kafka',
    'topic' = 'aiinfra.fct.userfeature.0',
    'properties.bootstrap.servers' = 'localhost:9092',
    'properties.group.id' = 'test-group',
    'format' = 'json'
)

使用 Flink,我执行 DDL。

val settings = EnvironmentSettings.newInstance.build
val streamEnv = StreamExecutionEnvironment.getExecutionEnvironment
val tableEnv = StreamTableEnvironment.create(streamEnv, settings)
tableEnv.executeSql(SOURCE_DDL)
val table = tableEnv.from("kafka_source")

然后,我把它转换成DataStream,在map(e => ...)部分做下游逻辑。

tableEnv.toRetractStream[(Long, java.sql.Timestamp, String)](table).map(e => ...)

map(e => ...) 部分,我想访问列名,在本例中为last_5_clicks。为什么?因为我可能有不同的来源,不同的列名(比如last_10min_page_view),但是我想复用map(e => ...)中的代码。

有没有办法做到这一点?谢谢。

【问题讨论】:

  • 这有什么意义??我的意思是在map 阶段您已经将数据转换为元组,因此您可以通过索引e._1 访问字段,因此您实际上不需要知道创建它们的列的名称。

标签: apache-flink flink-streaming flink-sql


【解决方案1】:

从 Flink 1.12 开始,它可以通过Table.getSchema.getFieldNames 访问。从1.13版本开始,可以通过Row.getFieldNames访问。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-03-22
    • 2018-09-12
    • 1970-01-01
    • 2018-12-09
    相关资源
    最近更新 更多