Spark SQL转Flink SQL:使用UNNEST创建临时表时test12未被识别求修复
Spark SQL转Flink SQL报错修复方案
原Spark SQL查询
CREATE OR REPLACE TEMPORARY VIEW tempView AS SELECT event.user_id AS userId, event.code AS code, test12.id AS id, test12 FROM user_interaction_data_table LATERAL VIEW explode(tests_explode) AS test12 WHERE event.user_id IS NOT NULL AND CAST(event.user_id AS INT) > 0 AND event.event_timestamp > 1710806400000 AND event.event_timestamp < 1710817200000;
尝试的Flink SQL查询
CREATE OR REPLACE TEMPORARY VIEW tempView AS SELECT event.user_id AS userId, event.code AS code, test12.id AS id, test12 FROM test_event UNNEST(tests_explode) AS test12 WHERE event.user_id IS NOT NULL AND CAST(event.user_id AS INT) > 0 AND event.event_timestamp > 1710806400000 AND event.event_timestamp < 1710817200000;
报错信息
test12is not recognized
备注信息
tests_explode是包含多个字段的嵌套结构体数组,结构如下:
tests_explode |--id |--field2
问题分析
Flink SQL中使用UNNEST展开数组时,必须结合LATERAL TABLE语法,通过别名表的形式明确指定展开后的元素别名,否则无法识别test12变量。
修复后的Flink SQL
CREATE OR REPLACE TEMPORARY VIEW tempView AS SELECT event.user_id AS userId, event.code AS code, test12.id AS id, test12 FROM test_event, LATERAL TABLE(UNNEST(tests_explode)) AS T(test12) WHERE event.user_id IS NOT NULL AND CAST(event.user_id AS INT) > 0 AND event.event_timestamp > 1710806400000 AND event.event_timestamp < 1710817200000;
验证说明
修复后的视图可正常执行后续查询:
SELECT test12.another_field from tempView;
内容的提问来源于stack exchange,提问作者stillLearning
相关产品推荐
相关产品推荐

