【问题标题】:How to complete aggregation task with flink cep如何使用 flink cep 完成聚合任务
【发布时间】:2019-10-14 16:05:33
【问题描述】:

我需要计算一天中发生 A 的次数和发生 B 的 15 分钟内的次数。流可能是 A1 ,A2,B1,B2,A3,B3,B4,B5,A4,A5,A6, A7,B6。 在我的情况下,事件结果是 A2,B1 A3,B3 A7,B6。当匹配器发生时,我需要接收实时结果。 我累了一些东西。我认为只有使用 flink cep 才能实现。但是 flink-sql-cep 不支持聚合。它只计算事件发生。在这种情况下,如何用一条 SQL 完成这个任务。

我累了两步。我先用flink sql cep来matcher,然后sink到kafka。在一步中,我使用 pre kafka 并使用 over window 进行聚合。

第一步: 选择引脚作为引脚,“第一步”作为 result_id,cast(order_amount as varchar) 作为 result_value,event_time 作为 result_time 来自 stra_dtpipeline MATCH_RECOGNIZE (按引脚分区
按 event_time 排序 措施
t1.pin 作为引脚, '1' 作为 order_amount, LOCALTIMESTAMP 作为 event_time 每场比赛一排 比赛后跳到下一行 间隔“30”秒内的模式 (t1 t2)
定义
t1 作为 t1.act_type='100001' , t2 作为 t2.act_type='100002' ) 第二步: select pin,'job5' as result_id,cast(sum(1) over (PARTITION BY pin,cast(DATE_FORMAT(event_time,'%Y%m%d') as VARCHAR) order by event_time ROWS BETWEEN INTERVAL '1' DAY PRECEDING AND CURRENT ROW ) as VARCHAR) as result_value,CURRENT_TIMESTAMP as result_time 来自 stra_dtpipeline_mid 其中 result_id='first-step' 和 DAYOFMONTH(CURRENT_DATE)=DAYOFMONTH(event_time)

我希望用一条 SQL 完成这项任务。

【问题讨论】:

    标签: apache-flink flink-cep flink-sql


    【解决方案1】:

    您可以使用子查询或视图将两个查询组合成一个查询。

    应该是这样的

    SELECT a, b OVER (...) ORDER BY event_time FROM (SELECT x, y MATCH_RECOGNIZE ...) WHERE ...
    

    CREATE VIEW pattern AS SELECT x, y MATCH_RECOGNIZE ...
    SELECT ... FROM pattern WHERE ...
    

    【讨论】:

      猜你喜欢
      • 2019-02-18
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多