【问题标题】:Spark Streaming job how to send data on Kafka topic and saving it in ElasticSpark Streaming 作业如何发送关于 Kafka 主题的数据并将其保存在 Elastic 中
【发布时间】:2019-06-04 20:16:40
【问题描述】:

我正在从事一个数据分析项目,在该项目中,我从 CSV 文件中读取数据,在 Kafka 主题上遍历该数据,并使用 Spark Streaming 来使用该 Kafka 主题数据。我在一个项目中使用的所有组件。

现在,在使用 Spark Streaming 消费数据之后,我必须对其进行一些计算,我必须将数据保存到弹性搜索中,并且我必须将这些数据发送到另一个主题。所以我正在从 Spark Streaming 做这些事情(将数据保存到弹性中并将数据发送到主题)。

下面是我的代码

@Component
public class RawEventSparkConsumer implements Serializable {

    @Autowired
    private ElasticSearchServiceImpl dataModelServiceImpl;

    @Autowired
    private EventKafkaProducer enrichEventKafkaProducer;

    Collection<String> topics = Arrays.asList("rawTopic");

    public void sparkRawEventConsumer(JavaStreamingContext streamingContext) {

        Map<String, Object> kafkaParams = new HashedMap();
        kafkaParams.put("bootstrap.servers", "localhost:9092");
        kafkaParams.put("key.deserializer", StringDeserializer.class);
        kafkaParams.put("value.deserializer", StringDeserializer.class);
        kafkaParams.put("group.id", "group1");
        kafkaParams.put("auto.offset.reset", "latest");
        kafkaParams.put("enable.auto.commit", true);

        JavaInputDStream<ConsumerRecord<String, String>> rawEventRDD = KafkaUtils.createDirectStream(streamingContext,
                LocationStrategies.PreferConsistent(),
                ConsumerStrategies.<String, String>Subscribe(topics, kafkaParams));

        JavaDStream<String> dStream = rawEventRDD.map((x) -> x.value());

        JavaDStream<BaseDataModel> baseDataModelDStream = dStream.map(convertIntoBaseModel);
        baseDataModelDStream.foreachRDD(rdd1 -> {
            saveDataToElasticSearch(rdd1.collect());
        });

        JavaDStream<EnrichEventDataModel> enrichEventRdd = baseDataModelDStream.map(convertIntoEnrichModel);

        enrichEventRdd.foreachRDD(rdd -> {
            System.out.println("Inside rawEventRDD.foreachRDD = = = " + rdd.count());
            sendEnrichEventToKafkaTopic(rdd.collect());
        });

        streamingContext.start();

        try {
            streamingContext.awaitTermination();
        } catch (InterruptedException e) {
            // TODO Auto-generated catch block
            e.printStackTrace();
        }

    }

    static Function convertIntoBaseModel = new Function<String, BaseDataModel>() {

        @Override
        public BaseDataModel call(String record) throws Exception {
            ObjectMapper mapper = new ObjectMapper();
            BaseDataModel csvDataModel = mapper.readValue(record, BaseDataModel.class);
            return csvDataModel;
        }
    };

    static Function convertIntoEnrichModel = new Function<BaseDataModel, EnrichEventDataModel>() {

        @Override
        public EnrichEventDataModel call(BaseDataModel csvDataModel) throws Exception {

            EnrichEventDataModel enrichEventDataModel = new EnrichEventDataModel(csvDataModel);
            enrichEventDataModel.setEnrichedUserName("Enriched User");
            User user = new User();
            user.setU_email("Nitin.Tyagi");
            enrichEventDataModel.setUser(user);
            return enrichEventDataModel;
        }
    };

    private void sendEnrichEventToKafkaTopic(List<EnrichEventDataModel> enrichEventDataModels) {
        if (enrichEventKafkaProducer != null && enrichEventDataModels != null && enrichEventDataModels.size() > 0)
            try {
                enrichEventKafkaProducer.sendEnrichEvent(enrichEventDataModels);
            } catch (JsonProcessingException e) {
                // TODO Auto-generated catch block
                e.printStackTrace();
            }
    }

    private void saveDataToElasticSearch(List<BaseDataModel> baseDataModelList) {
        if(!baseDataModelList.isEmpty())
            dataModelServiceImpl.saveAllBaseModel(baseDataModelList);
    }
}

现在我有几个问题

1) 我的方法是否可行,即将数据保存在 Elastic Search 中并从 Spark Streaming 发送到主题?

2) 我在单个项目中使用应用程序组件(Kafka、Spark Streaming),并且有多个 Spark Streaming 类。我正在本地系统中通过 CommandLineRunner 运行这些类。那么现在如何将 Spark Streaming 作为 Spark 作业提交呢?

对于 Spark Submit,我是否需要使用 Spark Streaming 类创建单独的项目?

【问题讨论】:

    标签: java spring-boot apache-spark apache-kafka spark-streaming


    【解决方案1】:

    我的方法是否可行,即在 Elastic Search 中保存数据并从 Spark Streaming 发送有关主题的数据?

    我想我会考虑使用 ES-Hadoop Spark 库。看起来您刚刚直接使用了 Elastic Java API(假设您正在收集 RDD 分区)

    虽然它可能有效,但它是高度耦合的......当 Elasticsearch 因维护而停机或高度潜在时会发生什么?整个应用程序是否停止?

    另一种方法是将 Kafka 处理逻辑拆分到自己的部署中。这样,您也可以只使用 Elasticsearch Kafka Connect 进程从主题加载数据,而无需自己编写该代码(Connect API 可能已经是您正在运行的 Kafka 集群的一部分)

    有多个 Spark Streaming 类

    多个主要课程?这不应该是一个问题。您需要为 Spark 提交提供一个 JAR 和一个类名。您可以在一个 jar 中拥有多个“入口点”/主要方法。

    如何将 Spark Streaming 作为 Spark 作业提交?

    我不确定我是否理解问题所在。 spark-submit 适用于流式作业

    注意:如果您计划更改数据类型或其顺序,CSV 是您可以放入 Kafka 的最差格式之一,并且您还希望除您自己以外的任何人都可以使用该主题。即使是 Elasticsearch 也希望你有 json 编码的有效负载

    【讨论】:

    • 谢谢@cricket_007,首先我会检查 Es-Spark 库。关于第二点,到目前为止,我已经在 Spring 中通过 CommandLineRunner 注册了 SparkStreaming 作业。那么在创建 jar 时,我应该通过 main 方法创建触发这些 Spark Streaming 作业吗?意味着我应该删除 commandlinerunner 并将 main() 方法放入 SparkStreamingJob。请指教。
    • 抱歉,不熟悉 Spring 如何与 Spark 协同工作
    • 好的。没问题,谢谢。你能给我一些示例链接吗,我们如何用 Java 编写和订阅 SparkStreaming?
    猜你喜欢
    • 2018-07-15
    • 2018-10-10
    • 2019-10-06
    • 2019-07-12
    • 2017-02-07
    • 2017-09-28
    • 2023-03-29
    • 2021-03-18
    • 1970-01-01
    相关资源
    最近更新 更多