【发布时间】:2019-01-16 02:10:36
【问题描述】:
我正在使用 Spark 应用程序来处理放置在我系统中 /home/user1/files/ 文件夹中的文本文件,并将这些文本文件中的逗号分隔数据映射为特定的 JSON 格式。我已经使用 spark 编写了以下 python 代码来做同样的事情。但是 Kafka 中的输出如下所示
Row(Name=Priyesh,Age=26,MailId=priyeshkaratha@gmail.com,Address=AddressTest,Phone=112)
Python 代码:
import findspark
findspark.init('/home/user1/spark')
from pyspark import SparkConf, SparkContext
from operator import add
import sys
from pyspark.streaming import StreamingContext
from pyspark.sql import Column, DataFrame, Row, SparkSession
from pyspark.streaming.kafka import KafkaUtils
import json
from kafka import SimpleProducer, KafkaClient
from kafka import KafkaProducer
producer = KafkaProducer(bootstrap_servers='server.kafka:9092')
def handler(message):
records = message.collect()
for record in records:
producer.send('spark.out', str(record))
print(record)
producer.flush()
def main():
sc = SparkContext(appName="PythonStreamingDirectKafkaWordCount")
ssc = StreamingContext(sc, 1)
lines = ssc.textFileStream('/home/user1/files/')
fields = lines.map(lambda l: l.split(","))
udr = fields.map(lambda p: Row(Name=p[0],Age=int(p[3].split('@')[0]),MailId=p[31],Address=p[29],Phone=p[46]))
udr.foreachRDD(handler)
ssc.start()
ssc.awaitTermination()
if __name__ == "__main__":
main()
那么如何在推送到 kafka 主题时将此行格式转换为 JSON?
【问题讨论】:
-
Spark 有 Kafka 库...为什么要收集 RDD 并使用常规的 Kafka 生产者?
-
@cricket_007 我是 Spark 和 Kafka 的新手。所以我一直在关注不同的教程,最终达到了这段代码。
-
好的,虽然这段代码可以将数据传输到 Kafka,但它并没有真正利用多台机器的 Spark 并行性
标签: apache-spark pyspark apache-kafka pyspark-sql