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:
- 导入Scala API包:
import org.apache.flink.streaming.api.scala._
- 修改ReducingStateDescriptor初始化代码:
val intervalDesc = new ReducingStateDescriptor[Long]("fire-time", Min(), createTypeInformation[Long]) val maxCountDesc= new ReducingStateDescriptor[Long]("count", Sum(), createTypeInformation[Long])
修改后报错即可消除。
内容的提问来源于stack exchange,提问作者Criwran
相关产品推荐
相关产品推荐

