【问题标题】:Send CSV from Kafka to Spark Streaming将 CSV 从 Kafka 发送到 Spark Streaming
【发布时间】:2017-09-28 19:29:07
【问题描述】:

我正在尝试将 csv 文件从 kafka 发送到 spark 流应用程序,但我不知道该怎么做。我在这里阅读了很多帖子,但没有人帮助我。

我希望我的 kafka 生产者发送 csv 并稍后在应用程序(消费者)中拆分它,但这并不重要。我试图创建一个 RDD 并将其发送给 spark。 这适用于普通字符串消息,但不适用于 csv

这是我的制作人:

message =sc.textFile("/home/guest/host/Seeds.csv")      
producer.send('test', message)

还有我的火花消费者:

ssc = StreamingContext(sc, 5)

kvs = KafkaUtils.createStream(ssc, "localhost:2181", "spark-streaming-consumer", {'test': 1}) data = kvs.map(lambda x: x[1]) counts = data.flatMap(lambda line: line.split(";")) \

.map(lambda word: (word, 1)) \
.reduceByKey(lambda a, b: a+b)

问题在于,通过发送 csv,火花流不会收到任何事件。 有人可以帮助我了解格式或概念吗?

我在 docker 容器下使用 python 在笔记本中运行生产者和消费者。

谢谢。

【问题讨论】:

    标签: python csv apache-spark streaming apache-kafka


    【解决方案1】:

    在我的工作中,我将任何 csv 转换为 json,

    这里是你如何在膝盖上做的例子(我的意思是没有任何import json

    from kafka import KafkaProducer
    import time,csv
    
    '''
    input csv example
    
    AT,BE,BG,CH,CY,CZ,DE,DK,EE,ES,FI,FR,EL,HR,HU,IE,IT,LT,LU,LV,NL,NO,PL,PT,RO,SI,SK,SE,UK
    0.15104895104895097,0.155978142670726,0.0,0.132959173102667,0,0.0261248185776488,0.0314454263905056,0.0,0.0,0.22130378970001602,0.0,0.0881265984931488,0.09026049932169501,0.056874262941565,0.0841602727424313,0.0494006197388216,0.0912473405767843,0.0,0.0656217442366246,0.0,0.0432966804004962,0.0,0.0,0.19138755980861197,0.0,0.0521335743946527,0.0,0.0,0.0434660616908725
    
    '''
    
    # create producer to kafka connection
    producer = KafkaProducer(bootstrap_servers='89.218.20.173:9092')
    # define *.csv file and a char that divide value
    fname = "input.csv"
    divider_char = ','
    # open file
    with open(fname) as fp:  
        # read header (first line of the input file)
        line = fp.readline()
        header = line.split(divider_char)
    
        #loop other data rows 
        line = fp.readline()    
        while line:
            # start to prepare data row to send
            data_to_send = ""
            values = line.split(divider_char)
            len_header = len(header)
            for i in range(len_header):
                data_to_send += "\""+header[i].strip()+"\""+":"+"\""+values[i].strip()+"\""
                if i<len_header-1 :
                    data_to_send += ","
            data_to_send = "{"+data_to_send+"}"
    
            '''
            example of outputs is valid JSON row 
            {
                "AT":"0.148251748251748",
                "BE":"0.052603706790461",
                    ...
                "SE":"0.0826699344612236",
                "UK":"0.10951678628072099"
            }
            '''
    
            # send data via producer
            producer.send('test', bytes(data_to_send, encoding='utf-8'))
            line = fp.readline()
            # А это так))) на всякий случай
            #time.sleep(1)
    producer.close()
    

    然后你可以使用下一个答案https://stackoverflow.com/a/47457985/6796393

    【讨论】:

      【解决方案2】:

      在您的生产者中,消息是一个 RDD(分布在集群中的 csv 文件行的集合),它被延迟评估,即在您对其执行操作之前它不会做任何事情。所以你需要在发送到 Kafka 之前收集 RDD。 请看下面的链接。 how to properly use pyspark to send data to kafka broker?

      【讨论】:

        猜你喜欢
        • 2017-02-07
        • 2020-08-04
        • 2015-11-26
        • 1970-01-01
        • 1970-01-01
        • 2019-08-08
        • 2016-04-23
        • 2017-12-06
        • 2016-03-12
        相关资源
        最近更新 更多