【发布时间】:2018-03-18 17:00:01
【问题描述】:
我正在尝试使用 scala 从火花流数据帧中提取值,其中包含如下代码:
var txs = spark.readStream
.format("kafka") .option("kafka.bootstrap.servers",KAFKABS)
.option("subscribe", "txs")
.load()
txs = txs.selectExpr("CAST(value AS STRING)")
val schema = StructType(Seq(
StructField("from",StringType,true),
StructField("to", StringType, true),
StructField("timestamp", TimestampType, true),
StructField("hash", StringType, true),
StructField("value", StringType, true)
))
txs = txs.selectExpr("cast (value as string) as json")
.select(from_json($"json", schema).as("data"))
.select("data.*")
.selectExpr("from","to","cast(timestamp as timestamp) as timestamp","hash","value")
val newDataFrame = txs
.flatMap(row => {
val to = row.getString(0)
val from = row.getString(1)
// val timestamp = row.getTimestamp??
//do stuff
})
我想知道时间戳是否有等效的类型化 get 方法?更让我感到困惑的是,我为结构化流定义的 SQL 类型与我通过flatMap 访问变量时的实际类型之间似乎存在某种隐藏映射(至少对我来说是隐藏的)功能。我查看了文档,确实是这样。根据文档:
返回位置 i 处的值。如果值为 null,则 null 为 回来。下面是 Spark SQL 类型和 返回类型:
BooleanType -> java.lang.Boolean ByteType -> java.lang.Byte
ShortType -> java.lang.Short IntegerType -> java.lang.Integer
FloatType -> java.lang.Float DoubleType -> java.lang.Double
StringType -> String DecimalType -> java.math.BigDecimalDateType -> java.sql.Date TimestampType -> java.sql.Timestamp
BinaryType -> 字节数组ArrayType -> scala.collection.Seq(使用 getList for java.util.List) MapType -> scala.collection.Map (使用 getJavaMap for java.util.Map) StructType -> org.apache.spark.sql.Row
考虑到这一点,我本来希望这个映射会更正式地作为它实现的接口被烘焙到 Row 类中,但显然情况并非如此 :( 似乎在 TimestampType 的情况下/java.sql.Timestamp,我必须放弃我的时间戳类型来做别的事情?有人请解释我为什么错了!我现在只使用 scala 和 spark 3-4 个月。
-保罗
【问题讨论】:
标签: scala apache-spark apache-spark-sql flatmap