You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink数据流增强:无合适顶层键时如何连接异吞吐量双数据流?

问题描述

我是Flink新手,需要实现有状态数据流增强:将Pizza Order流与更新频率更低的Pizza Price流连接,为订单中的每款披萨补充对应店铺的最新价格(取订单时间戳之前的最新价格)。尝试用keyBy算子时觉得两者没有合适的顶层键,附上我的Java实现代码,请教如何完成这个需求。

数据格式

1. Pizza Order

{
  "id": 123,
  "shop": "Mario's kitchen",
  "pizzas": [
    {
      "name": "Diavolo"
    },
    {
      "name": "Hawaii"
    }
  ],
  "timestamp": 12345678
}

2. Pizza Price

{
  "name": "Diavolo",
  "shop": "Mario's kitchen",
  "price": 14.2,
  "timestamp": 12345678
}

3. 期望输出:Enriched Pizza Order

{
  "id": 123,
  "shop": "Mario's kitchen",
  "pizzas": [
    {
      "name": "Diavolo",
      "price": 14.2
    },
    {
      "name": "Hawaii",
      "price": 12.5
    }
  ],
  "timestamp": 12345678
}

我的Java尝试代码

// 省略部分初始化代码
DataStream<PizzaOrder> pizzaOrderStream = env.fromSource(
        this.pizzaOrderSource,
        WatermarkStrategy.noWatermarks(),
        "Kafka Pizza Order Topic"
);

DataStream<PizzaPrice> pizzaPriceStream = env.fromSource(
        this.pizzaPriceSource,
        WatermarkStrategy.noWatermarks(),
        "Kafka Pizza Price Topic"
);

DataStream<EnrichedPizzaOrder> enrichedPizzaOrderDataStream =
        pizzaOrderStream.keyBy(PizzaOrder::getShop)
        .connect(pizzaPriceStream.keyBy(PizzaPrice::getShop))
        .process(new ProcessingTimeJoin());

enrichedPizzaOrderDataStream.sinkTo(sink);

return env.execute("Pizza Order Enriching Example");
}

public static class ProcessingTimeJoin extends CoProcessFunction<PizzaOrder, 
PizzaPrice, EnrichedPizzaOrder> {

        // 按店铺存储披萨名称与价格的映射
        private ValueState<HashMap<String, Double>> pizzaPriceState;

        @Override
        public void open(Configuration parameters) throws Exception {
            ValueStateDescriptor<HashMap<String, Double>> vDescriptor = new ValueStateDescriptor<>(
                    "pizzaPriceState",
                    TypeInformation.of(new TypeHint<HashMap<String, Double>>() {})
            );
            pizzaPriceState = getRuntimeContext().getState(vDescriptor);
        }

        @Override
        public void processElement1(PizzaOrder order,
                                    CoProcessFunction<PizzaOrder, PizzaPrice, EnrichedPizzaOrder>.Context ctx,
                                    Collector<EnrichedPizzaOrder> out) throws Exception {
            HashMap<String, Double> state = pizzaPriceState.value();
            if (state == null) {
                state = new HashMap<>();
            }
            List<EnrichedPizza> enrichedPizzas = new ArrayList<>();
            for (Pizza pizza : order.getPizzas()) {
                double price = state.getOrDefault(pizza.getPizzaName(), -1.0);

                EnrichedPizza newPizza = new EnrichedPizza(pizza, price);
                enrichedPizzas.add(newPizza);
            }
            EnrichedPizzaOrder enrichedPizzaOrder = new EnrichedPizzaOrder(order, enrichedPizzas);
            out.collect(enrichedPizzaOrder);
        }

        @Override
        public void processElement2(PizzaPrice price,
                                    CoProcessFunction<PizzaOrder, PizzaPrice, EnrichedPizzaOrder>.Context ctx,
                                    Collector<EnrichedPizzaOrder> out) throws Exception {
            HashMap<String, Double> state = pizzaPriceState.value();
            if (state == null) {
                state = new HashMap<>();
            }
            state.put(price.getName(), price.getPrice());
            pizzaPriceState.update(state);
        }
    }

解决方案

你的思路方向是对的,用shop作为keyBy的键完全可行——因为价格是和店铺绑定的,同一店铺的披萨价格更新和订单需要归到同一处理实例中。你的代码存在几个核心问题,调整后就能满足需求:

1. 切换到事件时间(Event Time)而非处理时间

你当前用的是WatermarkStrategy.noWatermarks(),这会导致Flink用处理时间计算,无法保证取到"订单时间戳之前的最新价格"。必须配置事件时间水位线:

步骤1:为两个数据流设置事件时间和水位线

// 为PizzaOrder设置事件时间与水位线
WatermarkStrategy<PizzaOrder> orderWatermark = WatermarkStrategy
        .<PizzaOrder>forMonotonousTimestamps()
        .withTimestampAssigner((order, timestamp) -> order.getTimestamp());

// 为PizzaPrice设置事件时间与水位线
WatermarkStrategy<PizzaPrice> priceWatermark = WatermarkStrategy
        .<PizzaPrice>forMonotonousTimestamps()
        .withTimestampAssigner((price, timestamp) -> price.getTimestamp());

DataStream<PizzaOrder> pizzaOrderStream = env.fromSource(
        this.pizzaOrderSource,
        orderWatermark,
        "Kafka Pizza Order Topic"
);

DataStream<PizzaPrice> pizzaPriceStream = env.fromSource(
        this.pizzaPriceSource,
        priceWatermark,
        "Kafka Pizza Price Topic"
);

如果数据流的时间戳不是单调递增的,可改用forBoundedOutOfOrderness(Duration.ofSeconds(5))来处理乱序数据,参数根据业务容忍的最大乱序时间调整。

2. 优化状态存储与价格更新逻辑

你的当前状态是用HashMap<String, Double>存储同一店铺的所有披萨价格,但无法保留价格的时间版本。如果需要严格保证"订单时间戳前的最新价格",应该为每个披萨存储带时间戳的价格历史,或者用MapState结合状态的TTL(Time-To-Live)来自动淘汰过期价格。

改进后的状态定义:用MapState存储单披萨的最新价格(带时间戳)

// 改为MapState:key是披萨名称,value是最新价格的时间戳与价格对
private MapState<String, Tuple2<Long, Double>> pizzaPriceState;

@Override
public void open(Configuration parameters) throws Exception {
    MapStateDescriptor<String, Tuple2<Long, Double>> descriptor = new MapStateDescriptor<>(
            "pizzaPriceState",
            BasicTypeInfo.STRING_TYPE_INFO,
            TypeInformation.of(new TypeHint<Tuple2<Long, Double>>() {})
    );
    // 可选:设置状态TTL,自动清理不再需要的旧价格(比如超过订单最大延迟时间的价格)
    StateTtlConfig ttlConfig = StateTtlConfig
            .newBuilder(Duration.ofHours(24))
            .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
            .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
            .build();
    descriptor.enableTimeToLive(ttlConfig);
    pizzaPriceState = getRuntimeContext().getMapState(descriptor);
}

更新价格的逻辑(processElement2)

@Override
public void processElement2(PizzaPrice price, Context ctx, Collector<EnrichedPizzaOrder> out) throws Exception {
    String pizzaName = price.getName();
    Long priceTimestamp = price.getTimestamp();
    Double priceValue = price.getPrice();

    // 仅当新价格的时间戳晚于已存储的最新价格时才更新
    Tuple2<Long, Double> existingPrice = pizzaPriceState.get(pizzaName);
    if (existingPrice == null || priceTimestamp > existingPrice.f0) {
        pizzaPriceState.put(pizzaName, Tuple2.of(priceTimestamp, priceValue));
    }
}

处理订单的逻辑(processElement1)

@Override
public void processElement1(PizzaOrder order, Context ctx, Collector<EnrichedPizzaOrder> out) throws Exception {
    Long orderTimestamp = order.getTimestamp();
    List<EnrichedPizza> enrichedPizzas = new ArrayList<>();

    for (Pizza pizza : order.getPizzas()) {
        String pizzaName = pizza.getPizzaName();
        Tuple2<Long, Double> priceEntry = pizzaPriceState.get(pizzaName);
        Double price = -1.0; // 默认值,代表无匹配价格

        // 仅取订单时间戳之前的最新价格
        if (priceEntry != null && priceEntry.f0 <= orderTimestamp) {
            price = priceEntry.f1;
        }

        enrichedPizzas.add(new EnrichedPizza(pizza, price));
    }

    out.collect(new EnrichedPizzaOrder(order, enrichedPizzas));
}

3. 处理无价格的场景

如果订单中的披萨还没有对应的价格更新,当前代码会用-1.0作为默认值。你可以根据业务需求调整:比如将这类订单暂存到状态中,等待价格更新后再输出,或者直接输出带默认值的订单。

如果需要暂存等待,可在processElement1中把未匹配到价格的订单存入ListState,然后在processElement2更新价格后,遍历状态中的订单重新检查并输出。


内容的提问来源于stack exchange,提问作者Mlemlord

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 07:55:02