insert into dim_result_lct_activy_config select Fact_id, LAST_VALUE(Fact_name), FROM_UNIXTIME(UNIX_TIMESTAMP(LAST_VALUE(Fact_start_time)), 'yyyyMMddHHmmss') as startTime, FROM_UNIXTIME(UNIX_TIMESTAMP(LAST_VALUE(Fact_end_time)), 'yyyyMMddHHmmss') as endTime, LAST_VALUE(Fstate) from db_act_config_t_act_logic_config groupby Fact_id
INSERT INTO mysql_sink SELECT f1, count(*) FROM kafka_src GROUP BY f1 每从 kafka 过来一条新的记录,会生成两条记录 Tuple2<Row, Boolean>, 旧的被删除,新的会添加上。这是query是会一个会产生retract stream的query,可以简单理解成每条kafka的数据过来会产生两条记录,但是最终写入下游的系统。需要看下游的系统支持和实现的sink(现在有三种sink AppendStreamSink, UpsertStreamSink, RetractStreamSink)
如若不带 group by 直接: INSERT INTO mysql_sink SELECT f1, f2 FROM kafka_src 主键冲突写入 mysql 是会出错的,怎么可以用 Upsert 的方式直接覆盖呢?
不带 group by时无法推导出query的 unique key,没法做按照unique key的更新, 只需要将 query的 key (你这里是group by 后的字段)和db中主键保持一致即可
TopN
1 2 3 4 5 6 7
SELECT [column_list] FROM ( SELECT [column_list], ROW_NUMBER() OVER ([PARTITIONBY col1[, col2...]] ORDERBY col1 [asc|desc][, col2 [asc|desc]...]) AS rownum FROM table_name) WHERE rownum <= N [AND conditions]
EXPLAIN PLAN FOR SELECT * FROM ( SELECT *, row_number() OVER(PARTITION BY merchandiseId ORDER BY totalQuantity DESC) AS rownum FROM ( SELECT merchandiseId, sum(quantity) AS totalQuantity FROM rtdw_dwd.kafka_order_done_log GROUP BY merchandiseId ) ) WHERE rownum <= 10
== Abstract Syntax Tree == LogicalProject(merchandiseId=[$0], totalQuantity=[$1], rownum=[$2]) +- LogicalFilter(condition=[<=($2, 10)]) +- LogicalProject(merchandiseId=[$0], totalQuantity=[$1], rownum=[ROW_NUMBER() OVER (PARTITION BY $0 ORDER BY $1 DESC NULLS LAST)]) +- ...
-- 滚动窗口统计pv:计算每分钟用户的页面访问数量 -- 上游表为tube,计算结果写入msql -- 假设上游表为页面点击记录表,计算pv只需要使用count函数进行条数统计即可 INSERT INTO console SELECT -- 滚动窗口的开始时间 FROM_UNIXTIME( CAST( TUMBLE_START(up.fevent_time, INTERVAL'1'MINUTE) asBIGINT ) /1000, 'yyyyMMddHHmmss' ) as start_time, count(1) as pv, count(distinct fuin) as uv FROM datagen AS up GROUPBY -- 按照1分钟进行开窗统计 -- 第一个参数指定事件时间 -- 第二个参数指定多久触发一次统计 TUMBLE(up.fevent_time, INTERVAL'1'MINUTE) -- 查看结果 -- 打开TM日志,全局搜索ConSoleSink output,即刻查看每分钟统计的数据
1 2 3 4 5 6 7 8 9 10 11 12
INSERT INTO t_hj_sample_result_040901 SELECT CAST(up.Fact_id ASINTEGER) AS catalog_01, CAST( TUMBLE_END(up.ptime, INTERVAL'10'MINUTE) ASBIGINT ) AS window_end FROM b_cdg_cft_lct_act_data_db_act_t_user_prize AS up GROUPBY CAST(up.Fact_id ASINTEGER), TUMBLE(up.ptime, INTERVAL'10'MINUTE)