【发布时间】:2020-10-13 03:02:26
【问题描述】:
我在定义的架构中有多种数据类型。试图找到一种按 TimestampType 过滤的好方法,将所有 TimeStampType 字段从 long 转换为 datetime。我可以在 StringType 的流中使用 .dtypes 进行过滤,但在尝试使用 StructFields 和 StructTypes 的 .dtypes 进行过滤时遇到问题。 有没有办法只过滤 Struct 中的 TimestampType? 下面是我在 Scala 2.11 中使用 spark 结构化流的 sudo 代码
val isoDateFormatter = "yyyy-MM-dd'T'HH:mm:ss'Z'"
val ExampleDataFrameLoad = spark
.readStream
.format("kafka")
.option("subscribe", topics.keys.mkString(","))
.options(kafkaConfig)
.load()
.select($"key".cast(StringType), $"value".cast(StringType), $"topic")
// Convert untyped dataframe to dataset
.as[(String, String, String)]
// Merge all manifests for vehicle in minibatch
.groupByKey(_._1)
//Start of merge
.flatMapGroupsWithState(OutputMode.Append, GroupStateTimeout.ProcessingTimeTimeout)(mergeGroup)
// .select($"key".cast(StringType),from_json($"value",schema).as("manifest"))
.select($"_1".alias("key"), $"_2".alias("jsonvalues"))
.select("key", "jsonvalues.*")
val ExampleDataFrame = ExampleDataFrameLoad
ExampleDataFrame.dtypes.foreach(println)
/* Returns
(key,StringType)
(contractVersion,StringType)
(metaData,StructType(StructField(Test,StringType,true), StructField(DateUtc,TimestampType,true)
*/
*Uses the following objects
import java.sql.Timestamp
object ManifestClasses {
final case class ProductManifestDocument(
contractVersion: Option[String],
metaData: DocumentMetaData
)
final case class DocumentMetaData(
Test: Option[String]
DateUtc: Timestamp
)
*/
ExampleDataFrame
//brings back data fields with types
.dtypes
//Currently returning empty but works for StringType
.filter(_._2 == "TimestampType")
.map(_._1)
//Tranforms all timestamp longs to yyyy-MM-dd'T'HH:mm:ss'Z' format
.foldLeft(ExampleDataFrame)((df, colName) => df.withColumn(colName, date_format(col(colName), isoDateFormatter)))
【问题讨论】:
标签: scala apache-spark spark-streaming