• flink cep数据源keyby union后 keybe失效


    问题背景:cep模板 对数据源设置分组条件后,告警的数据,和分组条件对不上, 掺杂了,其他的不同组的数据,产生了告警

    策略条件:

    选择了两个kafka的的topic的数据作为数据源,

    对A 数据源 test-topic1, 进行条件过滤, 过滤条件为:login_type  = 1

    对B 数据源 test-topic2,进行条件过滤,过滤条件为:login_type =  2

    分组条件 为   src_ip,hostname两个字段进行分组

    进行followby 关联。时间关联的最大时间间隔为  60秒

    运行并行度设置为3

    通过SourceStream打印的原始数据:

    1. 2> {"src_ip":"172.11.11.1","hostname":"hostname1","as":"A","create_time":1666859021060,"create_time_desc":"2022-10-27 16:23:41","event_type_value":"single","id":"67d32010-1f66-4850-b110-a7087e419c64_0","login_type":"1"}
    2. 2> {"src_ip":"172.11.11.1","hostname":"hostname1","as":"A","create_time":1666859020192,"create_time_desc":"2022-10-27 16:23:40","event_type_value":"single","id":"67d32010-1f66-4850-b110-a7087e419c64_0","login_type":"1"}
    3. 1> {"src_ip":"172.11.11.1","hostname":"hostname2","as":"B","create_time":1666859021231,"create_time_desc":"2022-10-27 16:23:41","event_type_value":"single","id":"67d32010-1f66-4850-b110-a7087e419c64_0","login_type ":"2"}

    经过cep处理后,产了告警

    产生告警:{A=[{"src_ip":"172.11.11.1","hostname":"hostname1","as":"A","create_time":1666859021060,"create_time_desc":"2022-10-27 16:23:41","event_type_value":"single","id":"67d32010-1f66-4850-b110-a7087e419c64_0","login_type":"1"}, {"src_ip":"172.11.11.1","hostname":"hostname1","as":"A","create_time":1666859020192,"create_time_desc":"2022-10-27 16:23:40","event_type_value":"single","id":"67d32010-1f66-4850-b110-a7087e419c64_0","login_type":"1"}], B=[{"src_ip":"172.11.11.1","hostname":"hostname2","as":"B","create_time":1666859021231,"create_time_desc":"2022-10-27 16:23:41","event_type_value":"single","id":"67d32010-1f66-4850-b110-a7087e419c64_0","login_type":"2"}]}
    

    经过src_ip,和hostname分组后, 理论上应该只分组后的相同的 scr_ip,hostname进行事件关联告警

    结果其他的分组数据也参和进来关联告警了。 

    期望的是  login_type = 1 出现至少两次, 接着login_type=2的至少出现1次,且相同的src_ip和hostname

    然后结果是下面数据也产生了告警。

    {"src_ip":"172.11.11.1","hostname":"hostname1","login_type":1}
    {"src_ip":"172.11.11.1","hostname":"hostname1","login_type":1}
    {"src_ip":"172.11.11.1","hostname":"hostname1","login_type":2}

    怀疑是分组没生效。

    然后debug数据源那块的方法kafkaStreamSource() 里面有进行分组,debug后发现确实也进行了keyby

    后来找不到其他问题,纠结了下, 怀疑是不是 KeyedSteam.union(KeyedStream)后得到的就不是一个KeyedSteam了。 所以

    出现问题的原始代码数据源代码:

    1. //程序具体执行流程
    2. DataStream sourceStream = SourceProcess.getKafkaStream(env, rule);
    3. DataStream resultStream = TransformProcess.process(sourceStream, rule);
    4. SinkProcess.sink(resultStream, rule);
    5. public static DataStream getKafkaStream(StreamExecutionEnvironment env, Rule rule) {
    6. DataStream inputStream = null;
    7. List events = rule.getEvents();
    8. if (events.size() > SharingConstant.NUMBER_ZERO) {
    9. for (Event event : events) {
    10. FlinkKafkaConsumer kafkaConsumer =
    11. new KafkaSourceFunction(rule, event).init();
    12. if (inputStream != null) {
    13. // 多条 stream 合成一条 stream
    14. inputStream = inputStream.union(kafkaStreamSource(env, event, rule, kafkaConsumer));
    15. } else {
    16. // 只有一条 stream
    17. inputStream = kafkaStreamSource(env, event, rule, kafkaConsumer);
    18. }
    19. }
    20. }
    21. return inputStream;
    22. }
    23. private static DataStream kafkaStreamSource(
    24. StreamExecutionEnvironment env,
    25. Event event,
    26. Rule rule,
    27. FlinkKafkaConsumer kafkaConsumer) {
    28. DataStream inputStream = env.addSource(kafkaConsumer);
    29. // 对多个黑白名单查询进行循环
    30. String conditions = event.getConditions();
    31. while (conditions.contains(SharingConstant.ARGS_NAME)) {
    32. // 使用新的redis 数据结构,进行 s.include 过滤
    33. inputStream = AsyncDataStream.orderedWait(inputStream,new RedisNameListFilterSourceFunction(s,rule.getSettings().getRedis()),30,TimeUnit.SECONDS,2000);
    34. conditions = conditions.replace(s, "");
    35. }
    36. // 一般过滤处理
    37. inputStream = AsyncDataStream.orderedWait(inputStream,
    38. new Redis3SourceFunction(event, rule.getSettings().getRedis()), 30, TimeUnit.SECONDS, 2000);
    39. // kafka source 进行 keyBy 处理
    40. return KeyedByStream.keyedBy(inputStream, rule.getGroupBy());
    41. }
    42. public static DataStream keyedBy(
    43. DataStream input, Map groupBy) {
    44. if (null == groupBy || groupBy.isEmpty() ||"".equals(groupBy.values().toArray()[SharingConstant.NUMBER_ZERO])){
    45. return input;
    46. }
    47. return input.keyBy(
    48. new TwoEventKeySelector(
    49. groupBy.values().toArray()[SharingConstant.NUMBER_ZERO].toString()));
    50. }
    51. public class TwoEventKeySelector implements KeySelector {
    52. private static final long serialVersionUID = 8534968406068735616L;
    53. private final String groupBy;
    54. public TwoEventKeySelector(String groupBy) {
    55. this.groupBy = groupBy;
    56. }
    57. @Override
    58. public String getKey(JSONObject event) {
    59. StringBuilder keys = new StringBuilder();
    60. for (String key : groupBy.split(SharingConstant.DELIMITER_COMMA)) {
    61. keys.append(event.getString(key));
    62. }
    63. return keys.toString();
    64. }
    65. }

    问题出现在这里:

    // 多条 stream 合成一条 stream
                        inputStream = inputStream.union(kafkaStreamSource(env, event, rule, kafkaConsumer));

    kafkaStreamSource()这个方法返回的是 KeyedStream ,

    两个KeyedStream unio合并后,  本来以为返回时KeyedStream,结果确是DataStream类型,

    结果导致cep分组不生效,一个告警中出现了其他分组的数据。

    解决方法, 就是在cep pattern前 根据是否有分组条件再KeyedBy一次

    1. private static DataStream patternProcess(DataStream inputStream, Rule rule) {
    2. PatternGen patternGenerator = new PatternGen(rule.getPatterns(), rule.getWindow().getSize());
    3. Pattern pattern = patternGenerator.getPattern();
    4. if (!rule.getGroupBy().isEmpty()){
    5. inputStream = KeyedByStream.keyedBy(inputStream, rule.getGroupBy());
    6. }
    7. PatternStream patternStream = CEP.pattern(inputStream, pattern);
    8. return patternStream.inProcessingTime().select(new RuleSelectFunction(rule.getAlarmInfo(), rule.getSelects()));

    输入数据:

    1.  {"src_ip":"172.11.11.1","hostname":"hostname1","as":"A","create_time":1666860300012,"create_time_desc":"2022-10-27 16:45:00","event_type_value":"single","id":"1288a709-d2b3-41c9-b7b7-e45149084514_0","login_type":"1"}
    2.  {"src_ip":"172.11.11.1","hostname":"hostname1","as":"A","create_time":1666860299272,"create_time_desc":"2022-10-27 16:44:59","event_type_value":"single","id":"1288a709-d2b3-41c9-b7b7-e45149084514_0","login_type":"1"}
    3.  {"src_ip":"172.11.11.1","hostname":"hostname2","as":"B","create_time":1666860300196,"create_time_desc":"2022-10-27 16:45:00","event_type_value":"single","id":"1288a709-d2b3-41c9-b7b7-e45149084514_0","login_type":"2"}

    不产生告警,符合预期

    再次输入同分组的数据:

    1. 2> {"src_ip":"172.11.11.1","hostname":"hostname1","as":"A","create_time":1666860369307,"create_time_desc":"2022-10-27 16:46:09","event_type_value":"single","id":"61004dd6-69ec-4d67-845c-8c15e7cc4bf7_0","app_id":"1"}
    2. 2> {"src_ip":"172.11.11.1","hostname":"hostname1","as":"A","create_time":1666860368471,"create_time_desc":"2022-10-27 16:46:08","event_type_value":"single","id":"61004dd6-69ec-4d67-845c-8c15e7cc4bf7_0","app_id":"1"}
    3. 2> {"src_ip":"172.11.11.1","hostname":"hostname1","as":"B","create_time":1666860369478,"create_time_desc":"2022-10-27 16:46:09","event_type_value":"single","id":"61004dd6-69ec-4d67-845c-8c15e7cc4bf7_0","app_id":"2"}
    4. 产生告警:{A=[{"src_ip":"172.11.11.1","hostname":"hostname1","as":"A","create_time":1666860368471,"create_time_desc":"2022-10-27 16:46:08","event_type_value":"single","id":"61004dd6-69ec-4d67-845c-8c15e7cc4bf7_0","app_id":"1"}, {"src_ip":"172.11.11.1","hostname":"hostname1","as":"A","create_time":1666860369307,"create_time_desc":"2022-10-27 16:46:09","event_type_value":"single","id":"61004dd6-69ec-4d67-845c-8c15e7cc4bf7_0","app_id":"1"}], B=[{"src_ip":"172.11.11.1","hostname":"hostname1","as":"B","create_time":1666860369478,"create_time_desc":"2022-10-27 16:46:09","event_type_value":"single","id":"61004dd6-69ec-4d67-845c-8c15e7cc4bf7_0","app_id":"2"}]}
    5. 告警输出:{"org_log_id":"61004dd6-69ec-4d67-845c-8c15e7cc4bf7_0,61004dd6-69ec-4d67-845c-8c15e7cc4bf7_0,61004dd6-69ec-4d67-845c-8c15e7cc4bf7_0","event_category_id":1,"event_technique_type":"无","event_description":"1","alarm_first_time":1666860368471,"src_ip":"172.11.11.1","hostname":"hostname1","intelligence_id":"","strategy_category_id":"422596451785379862","intelligence_type":"","id":"cc1cd8cd-a626-4916-bdd3-539ea57e898f","event_nums":3,"event_category_label":"资源开发","severity":"info","create_time":1666860369647,"strategy_category_name":"网络威胁分析","rule_name":"ceptest","risk_score":1,"data_center":"guo-sen","baseline":[],"sop_id":"","event_device_type":"无","rule_id":214,"policy_type":"pattern","strategy_category":"/NetThreatAnalysis","internal_event":"1","event_name":"ceptest","event_model_source":"/RuleEngine/OnLine","alarm_last_time":1666860369478}

    产生告警符合预期

  • 相关阅读:
    小程序--分包加载
    TensorBoard的使用1(add_scalar函数)
    计算机毕业设计Python+Django的学生作业管理系统
    C# .Net 发布后,把dll全部放在一个文件夹中,让软件目录更整洁
    真的够了,秋招已拿15+大厂offer,25k入职阿里全靠Java面试小册
    汽车电子——产品标准规范汇总和梳理(开发体系)
    精读DDD:service
    学习基因富集工具DAVID(3)
    时序预测|基于变分模态分解-时域卷积-双向长短期记忆-注意力机制多变量时间序列预测VMD-TCN-BiLSTM-Attention
    C++ STL 之顺序存储结构 vector,list,deque异同
  • 原文地址:https://blog.csdn.net/lr131425/article/details/127554972