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

Scala中Flink自定义Trigger报错:无法解析ReducingStateDescriptor构造方法

解决Flink自定义Trigger中ReducingStateDescriptor构造方法报错问题

问题场景

使用Scala 2.12.16搭配Flink 1.14.4自定义Trigger时,IDEA抛出以下错误:

Cannot resolve overloaded constructor ReducingStateDescriptor[Long]

报错代码片段:

val intervalDesc:ReducingStateDescriptor[Long] = new ReducingStateDescriptor[Long]("fire-time", Min(), LongSerializer.INSTANCE)
val maxCountDesc:ReducingStateDescriptor[Long] = new ReducingStateDescriptor[Long]("count", Sum(), LongSerializer.INSTANCE)

完整自定义Trigger代码:

import org.apache.flink.api.common.functions.ReduceFunction
import org.apache.flink.api.common.state.ReducingStateDescriptor
import org.apache.flink.api.common.typeutils.base.LongSerializer
import org.apache.flink.streaming.api.windowing.triggers.{Trigger, TriggerResult}
import org.apache.flink.streaming.api.windowing.windows.Window
import org.slf4j.{Logger, LoggerFactory}

class LqWaterLevelTrigger[W <: Window] extends Trigger[Object, W]{
  val log:Logger = LoggerFactory.getLogger(this.getClass)

  val interval:Long
  val maxCount:Long
  val intervalDesc:ReducingStateDescriptor[Long] = new ReducingStateDescriptor[Long]("fire-time", Min(), LongSerializer.INSTANCE)
  val maxCountDesc:ReducingStateDescriptor[Long] = new ReducingStateDescriptor[Long]("count", Sum(), LongSerializer.INSTANCE)


  override def onElement(element: Object, timestamp: Long, window: W, ctx: Trigger.TriggerContext): TriggerResult = ???

  override def onProcessingTime(time: Long, window: W, ctx: Trigger.TriggerContext): TriggerResult = ???

  override def onEventTime(time: Long, window: W, ctx: Trigger.TriggerContext): TriggerResult = ???

  override def clear(window: W, ctx: Trigger.TriggerContext): Unit = ???
}

case class Min() extends ReduceFunction[Long]  {
  override def reduce(value1: Long, value2: Long): Long = {
    Math.min(value1, value2)
  }
}

case class Sum() extends ReduceFunction[Long] {
  override def reduce(value1: Long, value2: Long): Long = {
    value1 + value2
  }
}

ReducingStateDescriptor的Java构造方法定义:

@PublicEvolving
public class ReducingStateDescriptor<T> extends StateDescriptor<ReducingState<T>, T> {

    private static final long serialVersionUID = 1L;

    private final ReduceFunction<T> reduceFunction;

    
    public ReducingStateDescriptor(
            String name, ReduceFunction<T> reduceFunction, Class<T> typeClass) {
        super(name, typeClass, null);
        this.reduceFunction = checkNotNull(reduceFunction);

        if (reduceFunction instanceof RichFunction) {
            throw new UnsupportedOperationException(
                    "ReduceFunction of ReducingState can not be a RichFunction.");
        }
    }

    
    public ReducingStateDescriptor(
            String name, ReduceFunction<T> reduceFunction, TypeInformation<T> typeInfo) {
        super(name, typeInfo, null);
        this.reduceFunction = checkNotNull(reduceFunction);
    }

    // 期望使用的构造方法
    public ReducingStateDescriptor(
            String name, ReduceFunction<T> reduceFunction, TypeSerializer<T> typeSerializer) {
        super(name, typeSerializer, null);
        this.reduceFunction = checkNotNull(reduceFunction);
    }

    
    public ReduceFunction<T> getReduceFunction() {
        return reduceFunction;
    }

    @Override
    public Type getType() {
        return Type.REDUCING;
    }
}

解决方案

添加Flink Scala API核心导入,并改用createTypeInformation[Long]替代LongSerializer.INSTANCE:

  1. 导入Scala API包:
import org.apache.flink.streaming.api.scala._
  1. 修改ReducingStateDescriptor初始化代码:
val intervalDesc = new ReducingStateDescriptor[Long]("fire-time", Min(), createTypeInformation[Long])
val maxCountDesc= new ReducingStateDescriptor[Long]("count", Sum(), createTypeInformation[Long])

修改后报错即可消除。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 22:57:24