【问题标题】:Spark Streaming collect()火花流收集()
【发布时间】:2018-03-15 20:05:47
【问题描述】:

我正在使用 Spark Streaming,并且我开发了以下 Spark Streaming 应用程序:

从 Kafka 接收器 (RDD1) 创建一个 DStream,从 HTTP 请求 (RDD2) 创建另一个。

我的问题是,我只想使用 RDD1 中的第一个元素并在我的 RDD2 中使用它,并且此代码在 spark 流式传输 (.first()) 中不起作用如何使用 spark 流式传输 1.6 获得相同的结果

代码:

   firstLineRDD = kvs.map(lambda x : x[0], x[1].split('\n')[0], x[2])
   dateRDD = firstLineRDD.map(lambda x : (datetime.datetime.fromtimestamp(float(x[0])/1000000),x[1],x[2]))
   dayAggRDD = dateRDD.map(lambda x : (x[0],x[1],x[2]))
   daily_date, sys , metric  = dayAggRDD.first()
   dataTSRDD = sc.parallelize(apiRequest(sys,metric,getDailyDate(daily_date)))

【问题讨论】:

  • RDD 还是DStreams ??
  • @massg 我已经使用 DStreams 修改了我的帖子
  • azelix,我还是没看到你指的DStreams

标签: apache-spark pyspark spark-streaming


【解决方案1】:

我在这里使用了两个 Kafka 流,但它与您的用例几乎相同。(在 java 中)。已将流转换为数据集,然后您可以将第一个元素用于第一个数据集,然后使用这些值进行操作第二个。以下是代码的简要摘要:

DstreamCoord.foreachRDD(new VoidFunction<JavaRDD<Coordinates>>() {

        String valueA;
        String valueB;

        public void call(JavaRDD<Coordinates> arg0) throws Exception {
            // TODO Auto-generated method stub

        RDD<Coordinates> d = arg0.rdd();

        Dataset<Row> dataset1 = session.createDataFrame(d, Coordinates.class);

        Row firstElement = dataset1.first();
            valueA = firstElement.getString(1);
            valueB = firstElement.getString(1);

            JavaInputDStream<ConsumerRecord<String, String>> stream2 =
                      org.apache.spark.streaming.kafka010.KafkaUtils.createDirectStream(
                        StreamingContext,
                        LocationStrategies.PreferConsistent(),
                        ConsumerStrategies.<String, String>Subscribe(topics, kafkaParams)
                      );

            stream2.foreachRDD(new VoidFunction<JavaRDD<ConsumerRecord<String,String>>>() {

                public void call(JavaRDD<ConsumerRecord<String, String>> arg0) throws Exception {
                    // TODO Auto-generated method stub

                Dataset<Row> dataset2 =     session.createDataFrame(arg0.rdd(), Coordinates.class);



                }
            });

【讨论】:

  • 这里我有两个数据集,其中有两列分别表示纬度和经度
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-04-27
  • 2020-11-02
  • 1970-01-01
  • 2019-04-02
  • 2016-02-07
  • 2015-05-15
相关资源
最近更新 更多