【问题标题】:Input data received all in lowercase on spark streaming in databricks using DataFrame使用 DataFrame 在数据块中的火花流上以小写形式接收的输入数据
【发布时间】:2019-11-26 17:37:30
【问题描述】:

我的 Spark 流应用程序使用来自 aws kenisis 的数据并部署在 databricks 中。我正在使用org.apache.spark.sql.Row.mkString 方法来使用数据,并且整个数据以小写形式接收。实际输入具有驼峰式字段名称和值,但在使用时以小写形式接收。

我尝试从一个简单的 java 应用程序中消费,并从 kinesis 队列中接收正确的数据。问题仅存在于使用 DataFrames 并在 databricks 中运行的 spark 流应用程序中。

// scala code

val query = dataFrame
      .selectExpr("lcase(CAST(data as STRING)) as krecord")
      .writeStream
      .foreach(new ForeachWriter[Row] {

        def open(partitionId: Long, version: Long): Boolean = {
          true
        }

        def process(row: Row) = {
          logger.info("Record received in data frame is -> " + row.mkString)
          processDFStreamData(row.mkString, outputHandler, kBase, ruleEvaluator)
        }

        def close(errorOrNull: Throwable): Unit = {
        }
      })
      .start()

期望火花流输入 json 应该是相同的情况 字母(驼峰式)作为 kinesis 中的数据,一旦使用数据帧接收,不应转换为小写。

有什么可能导致这种情况的想法吗?

【问题讨论】:

    标签: scala apache-spark dataframe databricks


    【解决方案1】:

    修复了这个问题,选择表达式中使用的 lcase 是罪魁祸首,更新了如下代码并且它工作了。

    val 查询 = 数据帧 .selectExpr("CAST(data as STRING) as krecord") .writeStream .foreach(new ForeachWriter[Row] { .........

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-07-07
      • 1970-01-01
      • 1970-01-01
      • 2012-08-18
      • 1970-01-01
      • 2016-09-14
      相关资源
      最近更新 更多