代码之家  ›  专栏  ›  技术社区  ›  Alejandro Alcalde

在Flink中调试自定义管道变压器

  •  1
  • Alejandro Alcalde  · 技术社区  · 8 年前

    我正试图在Flink中实现一个自定义转换器,如下所示 its documentation 但当我试图执行的时候 fit 从未调用操作。这就是我迄今为止所做的:

    class InfoGainTransformer extends Transformer[InfoGainTransformer] {
    
      import InfoGainTransformer._
    
      private[this] var counts: Option[collection.immutable.Vector[Map[Key, Double]]] = None
    
      // here setters for params, as Flink does
    
    }
    
    object InfoGainTransformer {
    
      // ====================================== Parameters =============================================
      // ...
    
      // ==================================== Factory methods ==========================================
      // ...
    
      // ========================================== Operations =========================================
    
      implicit def fitLabeledVectorInfoGain = new FitOperation[InfoGainTransformer, LabeledVector] {
        override def fit(instance: InfoGainTransformer, fitParameters: ParameterMap, input: DataSet[LabeledVector]): Unit = {
          val counts = collection.immutable.Vector[Map[Key, Double]]()
          input.map {
            v =>
              v.vector.map {
                case (i, value) =>
                  println("INSIDE!!!")
                  val key = Key(value, v.label)
                  val cval = counts(i).getOrElse(key, .0)
                  counts(i) + (key -> cval)
              }
          }
        }
      }
    
      implicit def fitVectorInfoGain[T <: Vector] = new FitOperation[InfoGainTransformer, T] {
        override def fit(instance: InfoGainTransformer, fitParameters: ParameterMap, input: DataSet[T]): Unit = {
          input
        }
      }
    
      implicit def transformLabeledVectorsInfoGain = {
        new TransformDataSetOperation[InfoGainTransformer, LabeledVector, LabeledVector] {
          override def transformDataSet(
                                         instance: InfoGainTransformer,
                                         transformParameters: ParameterMap,
                                         input: DataSet[LabeledVector]): DataSet[LabeledVector] = input
        }
      }
    
      implicit def transformVectorsInfoGain[T <: Vector : BreezeVectorConverter : TypeInformation : ClassTag] = {
        new TransformDataSetOperation[InfoGainTransformer, T, T] {
          override def transformDataSet(instance: InfoGainTransformer, transformParameters: ParameterMap, input: DataSet[T]): DataSet[T] = input
        }
      }
    }
    

    然后我试着用两种方法:

    val scaler = StandardScaler()
    val polyFeatures = PolynomialFeatures()
    val mlr = MultipleLinearRegression()
    val gain = InfoGainTransformer().setK(2)
    
    // Construct the pipeline
    val pipeline = scaler
      .chainTransformer(polyFeatures)
      .chainTransformer(gain)
      .chainPredictor(mlr)
    
    val r = pipeline.predict(dataSet map (_.vector))
    r.print()
    

    只有我的变压器:

    pipeline.fit(dataSet)
    

    在这两种情况下,当我在内部设置断点时 fitLabeledVectorInfoGain ,例如在行中 input.map 调试程序在那里停止,但是如果我在嵌套映射中也设置了断点,例如下面 println("INSIDE!!!") 它从不停在那里。

    有人知道如何调试这个自定义转换器吗?

    1 回复  |  直到 8 年前
        1
  •  0
  •   Alejandro Alcalde    8 年前

    FitOperation

    implicit def fitLabeledVectorInfoGain = new FitOperation[InfoGainTransformer, LabeledVector] {
        override def fit(instance: InfoGainTransformer, fitParameters: ParameterMap, input: DataSet[LabeledVector]): Unit = {
          //      val counts = collection.immutable.Vector[Map[Key, Double]]()
          val r = input.map {
            v =>
              v.vector.foldLeft(Map.empty[Key, Double]) {
                case (m, (i, value)) =>
                  println("INSIDE fit!!!")
                  val key = Key(value, v.label)
                  val cval = m.getOrElse(key, .0) + 1.0
                  m + (key -> cval)
              }
          }
          instance.counts = Some(r)
        }
      }
    

    TransformOperation

    推荐文章