【发布时间】:2021-11-18 13:46:51
【问题描述】:
我的任务是使用 structred spark streaming python 进行实时处理 所以第一步是将 csv 文件摄取到 kafka 主题中:完成
第二步是从主题kafka中读取流
df = spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", kafka_bootstrap_servers) \
.option("subscribe", kafka_topic_name) \
.option("startingOffsets", "latest")\
.load()
然后我使用模式来转换我的数据框的列
from pyspark.sql.types import *
schema = StructType() \
.add("DriverId",IntegerType(),True) \
.add("time",TimestampType(),True) \
.add("Longitude",DoubleType(),True) \
.add("Latitude",DoubleType(),True) \
.add("SPEED",DoubleType(),True) \
.add("EngineSpeed",IntegerType(),True) \
.add("MAF",IntegerType(),True) \
.add("FuelType",IntegerType(),True) \
第二步是在控制台上显示我正在播放的内容,看看我是否走对了:
query1 = df\
.writeStream\
.format("console")\
.outputMode("append")\
.option("truncate", False)\
.start()\
.awaitTermination()
所以,我已返回架构并将所有列转换为 stringType 结果是:
类似 json 的东西!
**我的问题是**如何正确转换我的数据框以及如何在每列下显示值和注释,如格式 json
【问题讨论】:
标签: csv apache-spark pyspark apache-kafka