【问题标题】:Why my MapState variable in Flink is not persisting previous values?为什么我在 Flink 中的 MapState 变量没有保留以前的值?
【发布时间】:2019-06-25 22:39:36
【问题描述】:

我正在用 Java 实现一个 Flink 程序来使用 MapStateDescriptor 处理状态。我基于此source 实现。出于某种原因,MapState 保留了以前的值,我无法计算出我想要的平均值。当我调试时,averageTemp 总是空的,我从来没有在里面找到任何密钥。我在实施过程中遗漏了什么?

import java.util.Map;

import org.apache.flink.api.common.functions.RichMapFunction;
import org.apache.flink.api.common.state.MapState;
import org.apache.flink.api.common.state.MapStateDescriptor;
import org.apache.flink.api.java.functions.KeySelector;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.TimeCharacteristic;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.sense.flink.mqtt.MqttTemperature;
import org.sense.flink.mqtt.TemperatureMqttConsumer;

public class SensorsMultipleReadingMqttEdgentQEP {

    public SensorsMultipleReadingMqttEdgentQEP() throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setStreamTimeCharacteristic(TimeCharacteristic.IngestionTime);

        DataStream<MqttTemperature> temperatureStream01 = env.addSource(new TemperatureMqttConsumer("topic-edgent-01"));
        DataStream<MqttTemperature> temperatureStream02 = env.addSource(new TemperatureMqttConsumer("topic-edgent-02"));
        DataStream<MqttTemperature> temperatureStream03 = env.addSource(new TemperatureMqttConsumer("topic-edgent-03"));
        DataStream<MqttTemperature> temperatureStreams = temperatureStream01.union(temperatureStream02)
                .union(temperatureStream03);

        DataStream<Tuple2<String, Double>> average = temperatureStreams.keyBy(new TemperatureKeySelector())
                .map(new AverageTempMapper());
        average.print();

        env.execute("SensorsMultipleReadingMqttEdgentQEP");
    }

    public static class TemperatureKeySelector implements KeySelector<MqttTemperature, Integer> {

        private static final long serialVersionUID = 5905504239899133953L;

        @Override
        public Integer getKey(MqttTemperature value) throws Exception {
            return value.getId();
        }
    }

    public static class AverageTempMapper extends RichMapFunction<MqttTemperature, Tuple2<String, Double>> {

        private static final long serialVersionUID = -5489672634096634902L;
        private MapState<String, Double> averageTemp;

        @Override
        public void open(Configuration parameters) throws Exception {
            averageTemp = getRuntimeContext()
                    .getMapState(new MapStateDescriptor<>("average-temperature", String.class, Double.class));
        }

        @Override
        public Tuple2<String, Double> map(MqttTemperature value) throws Exception {
            String key = "no-room";
            Double temp = value.getTemp();

            if (value.getId().equals(1) || value.getId().equals(2) || value.getId().equals(3)) {
                key = "room-A";
            } else if (value.getId().equals(4) || value.getId().equals(5) || value.getId().equals(6)) {
                key = "room-B";
            } else if (value.getId().equals(7) || value.getId().equals(8) || value.getId().equals(9)) {
                key = "room-C";
            }
            // NEVER ITERATES ON THE averageTemp
            for (Map.Entry<String, Double> entry: averageTemp.entries()) {
                System.out.println(entry.getKey() + " - " + entry.getValue());
            }

            System.out.println("value: " + value);
            if (averageTemp.contains(key)) { // NEVER CONTAINS A KEY
                System.out.println("yes: " + key);
                temp = (averageTemp.get(key) + value.getTemp()) / 2;
            } else {
                averageTemp.put(key, temp);
            }
            return Tuple2.of(key, temp);
        }
    }
}

**编辑:**好的。我误解了这个问题。该代码将先前的状态保存在 MapState 上。我错了,因为我正在调试代码。但实际上我遇到的问题是它启动了超过 1 个线程,并且在开始计算平均值之前它至少覆盖了我的地图值 3 次。

【问题讨论】:

    标签: java flink-streaming stateful


    【解决方案1】:

    您的地图函数中的状态是基于每个键。因此,当您的地图函数被调用并获得地图状态时,它将针对正在处理的MqttTemperature 记录中的任何 id。

    鉴于您想要每个房间的平均温度,我的处理方式如下:

    1. 根据 id 字段将TemperatureKeySelector 更改为返回room-Aroom-Broom-C
    2. AverageTempMapper 中,有两个ValueState 变量——一个是温度的总和(一个Double),另一个是一个计数。当你的map() 方法被调用时,如果这两个ValueState 变量中有一个为null,则将其初始化为0,然后求和/递增。

    【讨论】:

    • 不错。好办法。我要试试这个。谢谢
    • 嗨@kkrugler,我是这样做的(stackoverflow.com/questions/54535419/…)。你有更好的方法吗?
    • 您应该首先union() 三个流一起使用,然后将密钥提取器和平均应用于单个流。否则,每个流都会计算自己的平均值,这不是您想要的。
    • 是的。这种方式是一个更简单的查询计划。我试图在其中加入一些复杂性。
    猜你喜欢
    • 1970-01-01
    • 2019-12-13
    • 2019-07-07
    • 1970-01-01
    • 2022-11-27
    • 2020-12-09
    • 2017-01-28
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多