【问题标题】:Filter DataFrame based on words in array in Apache Spark根据 Apache Spark 中的数组中的单词过滤 DataFrame
【发布时间】:2019-03-07 21:13:12
【问题描述】:

我试图通过仅获取包含数组中单词的那些行来过滤数据集。 我正在使用 contains 方法,它适用于字符串但不适用于数组。下面是代码

val dataSet = spark.read.option("header","true").option("inferschema","true").json(path).na.drop.cache()

val threats_path = spark.read.textFile("src/main/resources/cyber_threats").collect()

val newData = dataSet.select("*").filter(col("_source.raw_text").contains(threats_path)).show()

它不起作用,因为威胁路径是字符串数组并且包含字符串的工作。任何帮助将不胜感激。

【问题讨论】:

  • 每一行是一个单词还是一个包含标点、空格等的短语?
  • 每一行有多个列。我试图过滤的列有很多文本、标点符号、空格等。

标签: scala apache-spark dataframe apache-spark-sql bigdata


【解决方案1】:

您可以在列上使用isinudf

它会变成这样,

val threats_path = spark.read.textFile("src/main/resources/cyber_threats").collect()

val dataSet = ???

dataSet.where(col("_source.raw_text").isin(thread_path: _*))

请注意,如果 thread_paths 的大小很大,这将影响性能,因为collect 和使用isin 的过滤器。

我建议您将过滤器dataSetthreats_path 结合使用join。它会像这样,

val dataSet = spark.read.option("header","true").option("inferschema","true").json(path).na.drop

val threats_path = spark.read.textFile("src/main/resources/cyber_threats")

val newData = threats_path.join(dataSet, col("_source.raw_text") === col("<col in threats_path >"), "leftouter").show()

希望对你有帮助

【讨论】:

  • 这里能不能写代码。在我可以申请加入的基础上,threshold_path 只有 1 列。
  • 是的,所以你用threats_path中的列名替换“”
  • 这些行正在显示,但我需要确认,它是否真的只显示了基于列 _source.raw_text 中含有位于威胁路径中的单词的行。尽管我使用了 na.drop,但仍然显示 null
  • 我要使用 isin。
  • threats_path 是文件,那么库名应该是什么?非?
猜你喜欢
  • 1970-01-01
  • 2016-02-05
  • 1970-01-01
  • 1970-01-01
  • 2016-12-28
  • 1970-01-01
  • 2017-02-14
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多