1 2 3 4 5 6 7 8
| Sink(table=[default_catalog.default_database.console1], fields=[user_id, data_date, data_time, msg_time, total_amount, video_exp_cnt]) +- Calc(select=[user_id, DATE_FORMAT(window_start, _UTF-16LE'yyyyMMdd') AS data_date, DATE_FORMAT(window_start, _UTF-16LE'HH:mm:ss') AS data_time, DATE_FORMAT(window_end, _UTF-16LE'yyyy-MM-dd HH:mm:ss') AS msg_time, total_amount, video_exp_cnt]) +- GlobalWindowAggregate(groupBy=[user_id], window=[CUMULATE(slice_end=[$slice_end], max_size=[86400000 ms], step=[10 min], offset=[0 ms], allowLazyTriggering=[true])], select=[user_id, SUM(sum$0) AS total_amount, COUNT(distinct$0 count$1) AS video_exp_cnt, start('w$) AS window_start, end('w$) AS window_end]) +- Exchange(distribution=[hash[user_id]]) +- LocalWindowAggregate(groupBy=[user_id], window=[CUMULATE(time_col=[order_time], max_size=[86400000 ms], step=[10 min], offset=[0 ms], allowLazyTriggering=[true])], select=[user_id, SUM(amount) AS sum$0, COUNT(distinct$0 order_id) AS count$1, DISTINCT(order_id) AS distinct$0, slice_end('w$) AS $slice_end]) +- Calc(select=[user_id, amount, order_id, order_time]) +- WatermarkAssigner(rowtime=[order_time], watermark=[-(order_time, 5000:INTERVAL SECOND)]) +- TableSourceScan(table=[[default_catalog, default_database, orders]], fields=[order_id, user_id, amount, order_time])
|