【问题标题】:Apache Flink - svm predictions on streaming dataApache Flink - 流数据的 svm 预测
【发布时间】:2019-10-23 12:09:54
【问题描述】:

我正在使用 Apache Flink 来预测来自 Twitter 的流。

代码在 Scala 中实现

我的问题是,我从 DataSet API 训练的 SVM-Model 需要一个 DataSet 作为 predict()-Method 的输入。

我已经在这里看到了一个问题,用户说,你需要编写一个自己的 MapFunction,它在工作开始时读取模型(参考:Real-Time streaming prediction in Flink using scala

但我无法编写/理解此代码。

即使我在 StreamingMapFunction 中获取模型。我仍然需要一个 DataSet 作为参数来预测结果。

我真的希望有人可以向我展示/解释这是如何完成的。

Flink 版本:1.9 Scala 版本:2.11 Flink-ML:2.11

val strEnv = StreamExecutionEnvironment.getExecutionEnvironment
val env = ExecutionEnvironment.getExecutionEnvironment

//this is my Model including all the terms to calculate the tfidf-values and to create a libsvm
val featureVectorService = new FeatureVectorService
        featureVectorService.learnTrainingData(labeledData, false)

//reads the created libsvm
val trainingData: DataSet[LabeledVector] = MLUtils.readLibSVM(env, "...")
        val svm = SVM()
                .setBlocks(env.getParallelism)
                .setIterations(100)
                .setRegularization(0.001)
                .setStepsize(0.1)
                .setSeed(42)
//learning
svm.fit(trainingData)

//this is my twitter stream - text should be predicted later
val streamSource: DataStream[String] = strEnv.addSource(new TwitterSource(params.getProperties))

//the texts i want to transform to tfidf using the service upon and give it the svm to predict
val tweets: DataStream[(String, String)] = streamSource
                .flatMap(new SelectEnglishTweetWithCreatedAtFlatMapper)

【问题讨论】:

  • 你想用 java 或 scala 做吗?
  • @AntonioMiranda 嘿,感谢您的回答。我在斯卡拉做。更新了上面的代码以便更好地理解。希望你能帮忙

标签: dataset apache-flink data-stream flinkml


【解决方案1】:

因此,SVM 所属的 FlinkML 目前不支持流 API。这就是为什么SVM 只接受DataSet。这个想法不是使用 FlinkML,而是使用 scala 或 java 中可用的一些 SVM 库。然后您可以读取模型,例如从文件中读取。问题是您必须自己实现大部分逻辑。

你提到的帖子中的评论或多或少说的是完全相同的事情。

【讨论】:

  • 感谢您的回答 :) 我知道 FlinkML 的 SVM 本身不支持流式传输。但是我发现我链接了这篇文章,它看起来可能带有一些提示。只是我不明白它背后的逻辑。
  • 是的,您提到的帖子的答案与我所描述的差不多:) 您应该从文件中读取模型并存储为MapFunction 中的状态.里面提到的Model类型只是一个接口,不多说:ci.apache.org/projects/flink/flink-docs-release-1.9/api/java/…
  • 从文件中读取模型意味着在这种情况下我的 SVM?但是我如何让它接受一个 Vector 或 DataString[Vector] 来预测? (再次感谢您的回答)
  • 您可能会使用一些可用于 Java 的 SVM 库,例如 csie.ntu.edu.tw/~cjlin/libsvm 或任何其他库 :) 但您肯定需要对其进行调整以允许使用 Vector
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-07-15
  • 2017-01-05
  • 1970-01-01
  • 1970-01-01
  • 2019-06-21
  • 2019-02-15
相关资源
最近更新 更多