【问题标题】:SplitStream for dynamic output key (Select)SplitStream 用于动态输出键(选择)
【发布时间】:2017-05-18 14:19:36
【问题描述】:

这是我的代码。

    SplitStream<MonitoringEvent> splitStream =  inputStream.split(new OutputSelector<MonitoringEvent>() {

    @Override
    public Iterable<String> select(MonitoringEvent me) {

        List<String> ml = new ArrayList<String>();              
        ml.add(me.getEventType());                              
        return ml;
}

我有随机顺序的监控事件流 温度:80,压力:70,湿度:80,温度:30...

使用上面的代码,我正在拆分流,事件类型,即温度流,压力流。

问题是,如果我知道 eventType,我可以像这样从 splitStream 中选择它

splitStream.select('temperatureStream')

但 eventType 是动态的,不是预定义的。

我将如何为这个动态流应用 CEP。如果

temperate is > 90 for past 10 minutes ...

pressure is > 90 for past 10 minutes ...

【问题讨论】:

  • 不理想,但由于您的事件类型是有限且小的(温度、压力、湿度...),您可以拥有多个流,然后对这些单独的流进行类型特定的处理。如果 eventTypes 显着增长,那么是的,这将很难管理。
  • 或在源/生产者或使用某种基于键的路由(如消息传递)预先拆分事件
  • @madhairsilence 这个有什么解决方案吗?我几乎有同样的问题。
  • 您不能为此使用 Flink 的 CEP 组件。您将不得不编写自定义窗口事件。并处理它。如果可能,将尝试发布代码。

标签: apache-flink complex-event-processing data-stream


【解决方案1】:

如果我错了,请纠正我,但我认为不可能对 select due flink 的并行性进行动态查找。你的程序被翻译成 flinks taskmanagers 的并行指令,jobmanager 协调这些动作。如果没有对抽象语法树的全面了解,根本就无法应用并行性......也许你可以找到所有消息共享和不同的一些共同属性

【讨论】:

  • 那么我该如何实现呢?具有多种事件类型的单个流。并且必须相应地拆分流并应用 CEP
  • 如果您还没有找到解决方案。您必须使用源上的 DeserializationSchema 指定传入消息并将其解析为类。这里有一些伪代码,例如:class mySchema extends AbstractDeserializationSchema { public LogEvent deserialize(byte[] message) throws IOException {return (Temperature)message;} } 每个传入的消息都被键入为温度类,并且可以选择每个通过使用 .select().type(Temperature.class).dosomething 和 CEP-Rules 的事件
猜你喜欢
  • 2019-06-10
  • 1970-01-01
  • 2021-01-31
  • 2017-03-05
  • 2016-10-20
  • 2016-04-01
  • 2015-10-04
  • 2019-03-03
  • 1970-01-01
相关资源
最近更新 更多