【问题标题】:Flink DataStream map Object to list of ObjectsFlink DataStream 将对象映射到对象列表
【发布时间】:2021-11-19 05:46:09
【问题描述】:

我正在尝试将对象 A 的 DataStream 转换为对象 B 的列表。如下例所示,我正在从 flink 消费者读取 DataStream,我需要转换为 DataStream 以便我可以运行一些过滤器和聚合MappedMetric 对象上的 timeWindow。一个 LogEvent 可能会导致 MappedMetric 对象列表,因此如果我使用 MapFunction,结果将是 DataStream。但是,我认为聚合不能在 DataStream 上运行。非常感谢任何帮助。提前致谢。

// Input Object
public class LogEvent {
    private String id;
    private long timestamp;
    private List<LogMessage> message;
}

public class LogMessage {
    private String accountId;
    private List<Metric> metrics;
}

public class Metric {
    private String name;
    private double value;
}

// Should be transformed to 

public class MappedMetric {
    private String accountId;
    private String name;
    private double value;
    private long timestamp;
}

final DataStream<LogEvent> inputDataStream = **read from Flink consumer**
final DataStream<MappedMetric> aggregatedMetrics = inputDataStream
                .map(**SomeMapFunction**)
                .keyBy(**SomeKey**)
                
                
                

【问题讨论】:

  • 您能否更好地解释您打算如何处理您的数据?您有 LogEvent 输入流,您想将它与 LogMessage 和 Metric 一起加入吗?那些也是流吗?您想与ListState<T> 合作吗?你能提供一个数据输入和预期输出的例子吗?

标签: java apache-flink flink-streaming


【解决方案1】:

您想使用FlatMap 函数,它可以为单个输入生成多个结果。每个结果都是一个 MappedMetric 记录,而不是一个列表。

【讨论】:

    猜你喜欢
    • 2013-03-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-12-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-11-16
    相关资源
    最近更新 更多