【发布时间】: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