【发布时间】: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