【问题标题】:Spark streaming duplicate network calls火花流重复网络调用
【发布时间】:2017-04-06 18:00:43
【问题描述】:

我正在使用 pyspark 和 Kafka 接收器来处理推文流。我的应用程序的其中一个步骤包括调用 Google Natural Language API 以获取每条推文的情绪分数。但是,我看到 API 在每条已处理的推文中都会收到多次调用(我在 Google Cloud Console 中看到了调用次数)。

另外,如果我打印 tweetID(在映射函数内),我会得到相同的 ID 3 或 4 次。在我的应用程序结束时,推文被发送到 Kafka 中的另一个主题,在那里我得到了正确的推文计数(没有重复的 ID),所以原则上一切正常,但我不知道如何避免调用谷歌每条推文不止一次 API。

这是否与 Spark 或 Kafka 中的某些配置参数有关?

这是我的控制台输出示例:

TIME 21:53:36: Google Response for tweet 801181843500466177 DONE!
TIME 21:53:36: Google Response for tweet 801181854766399489 DONE!
TIME 21:53:36: Google Response for tweet 801181844808966144 DONE!
TIME 21:53:37: Google Response for tweet 801181854372012032 DONE!
TIME 21:53:37: Google Response for tweet 801181843500466177 DONE!
TIME 21:53:37: Google Response for tweet 801181854766399489 DONE!
TIME 21:53:37: Google Response for tweet 801181844808966144 DONE!
TIME 21:53:37: Google Response for tweet 801181854372012032 DONE!

但在 Kafka 接收器中,我只收到 4 条已处理的推文(这是正确的接收方式,因为它们只有 4 条独特的推文)。

执行此操作的代码是:

def sendToKafka(rdd,topic,address):
    publish_producer = KafkaProducer(bootstrap_servers=address,\
                            value_serializer=lambda v: json.dumps(v).encode('utf-8'))
    records = rdd.collect()
    msg_dict = defaultdict(list)
    for rec in records:
        msg_dict["results"].append(rec)
    publish_producer.send(resultTopic,msg_dict)
    publish_producer.close()


kafka_stream = KafkaUtils.createStream(ssc, zookeeperAddress, "spark-consumer-"+myTopic, {myTopic: 1})

dstream_tweets=kafka_stream.map(lambda kafka_rec: get_json(kafka_rec[1]))\
                 .map(lambda post: add_normalized_text(post))\
                 .map(lambda post: tagKeywords(post,tokenizer,desired_keywords))\
                 .filter(lambda post: post["keywords"] == True)\
                 .map(lambda post: googleNLP.complementTweetFeatures(post,job_id))

dstream_tweets.foreachRDD(lambda rdd: sendToKafka(rdd,resultTopic,PRODUCER_ADDRESS))

【问题讨论】:

  • 你已经做了什么?您能否将您的代码粘贴到问题中?
  • 我用代码更新了问题。 googleNLP.complementTweetFeatures() 向 Google API 发出一个请求并返回响应。

标签: apache-spark apache-kafka pyspark spark-streaming twitter-streaming-api


【解决方案1】:

我已经找到了解决方案!我只需要缓存 DStream:

dstream_tweets.cache()

发生多个网络调用是因为 Spark 在我的脚本中执行后面的操作之前重新计算了该 DStream 中的 RDD。当我缓存() DStream 时,只需要计算一次;并且由于它保存在内存中,以后的函数可以访问该信息而无需重新计算(在这种情况下涉及重新计算以再次调用 API,因此付出更多内存使用的代价是值得的)。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-04-27
    • 2018-08-09
    • 2019-04-02
    • 2016-02-07
    • 2015-05-15
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多