【发布时间】:2020-11-27 03:41:11
【问题描述】:
遇到一个使用翻转窗口的 apache flink 应用程序的问题。窗口大小为 10 秒,我希望每 10 秒有一个 resultSet DataStream。但是,当最新窗口的结果集总是延迟时,除非我将更多数据推送到源流。
例如,如果我在 '01:33:40.0' 和 '01:34:00.0' 之间将多条记录推送到源流,然后停下来查看日志,什么都不会发生。
我再次在“01:37:XX”上推送一些数据,然后将获得“01:33:40.0”和“01:34:00.0”之间窗口的结果集,这是不期望的,因为下游接收器逻辑期待结果集准时。
非常感谢任何改进这一点的提示。谢谢。
以下是日志:
"log timestamp": "2019-11-15 01:37:45",
"message": "resultSet output: CLASS: 13 CNT: 1 from: 2019-11-15 01:33:40.0 to: 2019-11-15 01:34:00.0\n",
下面是代码sn-p:
Table resultTable = tableEnv.sqlQuery(""+
"SELECT " +
" CAST (N02_001 AS VARCHAR(10)) AS RAILWAY_CLASS, " +
" COUNT(*) RAILWAY_CLASS_COUNT, " +
" TUMBLE_START(rowtime, INTERVAL '20' SECOND) as WINDOW_START, " +
" TUMBLE_END(rowtime, INTERVAL '20' SECOND) as WINDOW_END " +
" FROM Inputs " +
" GROUP BY TUMBLE(rowtime, INTERVAL '20' SECOND), CAST (N02_001 AS VARCHAR(10))");
TupleTypeInfo<Tuple4<String, Long, Timestamp, Timestamp>> tupleType = new TupleTypeInfo<>(
Types.STRING,
Types.LONG,
Types.SQL_TIMESTAMP,
Types.SQL_TIMESTAMP);
DataStream<Tuple4<String, Long, Timestamp, Timestamp>> resultSet = tableEnv.toAppendStream(resultTable, tupleType);
resultSet
.map((Tuple4<String, Long, Timestamp, Timestamp> value) -> {
String output = "CLASS: " + value.f0 + " CNT: " + value.f1 + " from: " + value.f2 + " to: " + value.f3 + "\n";
log.warn("resultSet output: " + output);
return value;
})
.returns(Types.TUPLE(Types.STRING, Types.LONG, Types.SQL_TIMESTAMP, Types.SQL_TIMESTAMP));
【问题讨论】:
标签: apache-flink