【问题标题】:flink scala map with dead letter queue带有死信队列的flink scala映射
【发布时间】:2022-01-19 09:49:11
【问题描述】:

我正在尝试制作一些 scala 函数来帮助进行 flink mapfilter 操作,将它们的错误重定向到死信队列。

但是,我正在与 scala 的类型擦除作斗争,这使我无法将它们设为通用。下面mapWithDeadLetterQueue的实现不编译。


sealed trait ProcessingResult[T]
case class ProcessingSuccess[T,U](result: U) extends ProcessingResult[T]
case class ProcessingError[T: TypeInformation](errorMessage: String, exceptionClass: String, stackTrace: String, sourceMessage: T) extends ProcessingResult[T]

object FlinkUtils {
    // https://stackoverflow.com/questions/1803036/how-to-write-asinstanceofoption-in-scala
    implicit class Castable(val obj: AnyRef) extends AnyVal {
        def asInstanceOfOpt[T <: AnyRef : ClassTag] = {
            obj match {
            case t: T => Some(t)
            case _ => None
            }
        }
    }

    def mapWithDeadLetterQueue[T: TypeInformation,U: TypeInformation](source: DataStream[T], func: (T => U)): (DataStream[U], DataStream[ProcessingError[T]]) = {
        val mapped = source.map(x => { 
            val result = Try(func(x)) 
            result match {
                case Success(value) => ProcessingSuccess(value)
                case Failure(exception) => ProcessingError(exception.getMessage, exception.getClass.getName, exception.getStackTrace.mkString("\n"), x)
            }
        } )
        val mappedSuccess = mapped.flatMap((x: ProcessingResult[T]) => x.asInstanceOfOpt[ProcessingSuccess[T,U]]).map(x => x.result)
        val mappedFailure = mapped.flatMap((x: ProcessingResult[T]) => x.asInstanceOfOpt[ProcessingError[T]])
        (mappedSuccess, mappedFailure)
    }
  
}

我明白了:

[error] FlinkUtils.scala:35:36: overloaded method value flatMap with alternatives:
[error]   [R](x$1: org.apache.flink.api.common.functions.FlatMapFunction[Product with Serializable with ProcessingResult[_ <: T],R], x$2: org.apache.flink.api.common.typeinfo.TypeInformation[R])org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator[R] <and>
[error]   [R](x$1: org.apache.flink.api.common.functions.FlatMapFunction[Product with Serializable with ProcessingResult[_ <: T],R])org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator[R]
[error]  cannot be applied to (ProcessingResult[T] => Option[ProcessingSuccess[T,U]])
[error]         val mappedSuccess = mapped.flatMap((x: ProcessingResult[T]) => x.asInstanceOfOpt[ProcessingSuccess[T,U]]).map(x => x.result)

有没有办法让它工作?

【问题讨论】:

    标签: scala apache-flink flink-streaming


    【解决方案1】:

    好的,我要回答我自己的问题。我犯了几个错误:

    • 首先,我不小心包含了 java DataStream 类而不是 scala DataStream 类(这种情况经常发生)。 java 变体显然不接受 map/filter/flatmap 的 scala lambda
    • 其次,flink 序列化不支持密封特征。 a project 应该可以解决它,但我还没有尝试过。

    解决方案:首先我没有使用密封特性,而是使用带有两个选项的简单案例类(表达力稍差,但仍然有效):

    case class ProcessingError[T](errorMessage: String, exceptionClass: String, stackTrace: String, sourceMessage: T)
    case class ProcessingResult[T: TypeInformation, U: TypeInformation](result: Option[U], error: Option[ProcessingError[T]])
    

    然后,我可以让一切都像这样工作:

    object FlinkUtils {
        def mapWithDeadLetterQueue[T: TypeInformation: ClassTag,U: TypeInformation: ClassTag]
           (source: DataStream[T], func: (T => U)): 
           (DataStream[U], DataStream[ProcessingError[T]]) = {
            implicit val typeInfo = TypeInformation.of(classOf[ProcessingResult[T,U]])
    
            val mapped = source.map((x: T) => { 
                val result = Try(func(x)) 
                result match {
                    case Success(value) => ProcessingResult[T, U](Some(value), None)
                    case Failure(exception) => ProcessingResult[T, U](None, Some(
                      ProcessingError(exception.getMessage, exception.getClass.getName, 
                               exception.getStackTrace.mkString("\n"), x)))
                }
            } )
            val mappedSuccess = mapped.flatMap((x: ProcessingResult[T,U]) => x.result)
            val mappedFailure = mapped.flatMap((x: ProcessingResult[T,U]) => x.error)
            (mappedSuccess, mappedFailure)
        }
    
    

    }

    flatMapfilter 函数看起来非常相似,但它们分别使用 ProcessingResult[T,List[T]]ProcessingResult[T,T]

    我使用这样的函数:

    val (result, errors) = FlinkUtils.filterWithDeadLetterQueue(input, (x: MyMessage) => {
              x.`type` match {
                case "something" => throw new Exception("how how how")
                case "something else" => false
                case _ => true
              }
    })
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2018-10-28
      • 1970-01-01
      • 2012-05-02
      • 2018-04-16
      • 2021-10-07
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多