记录首次使用kafka消费数据到hive:

服务端spark版本:2.1 kafka:2.0.0

1.开始用2.1版本的spark客户端依赖,该版本貌似只适合用spark streaming 方式,使用structured streaming兼容性不好,整半天没行,该版本自带的kafka客户端是0.1.0(大坑!!!貌似不能设置MAX_POLL_INTERVAL_MS_CONFIG参数),正常写入但是提交偏移量的时候一直报错再平衡,我调了好多次都没有用(我是先使用本地调试)。

<dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-streaming-kafka-0-10_2.11</artifactId>
            <version>2.1.0</version>
</dependency>

2.后续手动排除自带的kafka版本,我看官方文档其实不建议这么做他说集成的kafka版本强依赖,自己改了可能出现问题;但是我试了这样确实也可以,但是我只在设置自动提交偏移量的模式成功,手动模式一直没成功也是再平衡问题,主要就是这个偏移量问题。

<dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-streaming-kafka-0-10_2.11</artifactId>
            <version>2.1.0</version>
           <exclusions>
                <exclusion>
                   <groupId>org.apache.kafka</groupId>
                   <artifactId>kafka-clients</artifactId>
                </exclusion>
           </exclusions>
</dependency>


<dependency>
       <groupId>org.apache.kafka</groupId>
       <artifactId>kafka-clients</artifactId>
       <version>2.0.0</version>
</dependency>

3.最后在我几十次调试之后,直接上了2.4.8的客户端版本,这个版本我之前hive-redis的时候用过,没遇到问题(实在是我们生产2.1.0太老了),这个2.4.8版本自带的是kafka 2.0.0客户端,试了一下居然可以了,想哭了之前没用这个版本试了n次都不行。

<dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-streaming-kafka-0-10_2.11</artifactId>
            <version>2.4.8</version>
 </dependency>

最后我的kafka设置如下:

// Kafka配置参数
        Map<String, Object> kafkaParams = new HashMap<>();
        kafkaParams.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaBootstrapServers);
        kafkaParams.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
        kafkaParams.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
        kafkaParams.put(ConsumerConfig.GROUP_ID_CONFIG, consumerGroupId);
        kafkaParams.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); // 首次运行从最早开始
        kafkaParams.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
        //kafkaParams.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, "2000"); // 2秒提交一次
        kafkaParams.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, "600000"); // 10分钟
        //kafkaParams.put("max.poll.interval.ms", "600000"); // 10分钟
        kafkaParams.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "45000"); // 会话超时设为 45 秒
        kafkaParams.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, "10000"); // 心跳间隔设为 10 秒

另外还有一些细节问题:

1.stream流消费者对象不能直接在rdd后面的处理里面直接使用,例如用它提交偏移量会直接报错不能序列化,可以先像下面一样强转一下,获取偏移量也必须在最开始的rdd对象中获取。可能是因为这些版本不兼容,我开始这么用也没行,主要就是再平衡导致失败或者不能序列化,后面我单独创建了一个临时消费者用于提交偏移量,解决了不能序列化问题,但是又是再平衡导致失败。

CanCommitOffsets canCommitOffsets = (CanCommitOffsets) stream.inputDStream();
// 或者直接
((CanCommitOffsets) stream.inputDStream()).commitAsync()

2.其实最开始我没有用kafka自带的偏移量提交方式,我使用的是checkpoint,目录选择的我自己的hdfs目录,但是没有成功,检查点目录下正常写入数据后,我直接kill -9 终止程序,然后重启 没办法正常从检查点拿到偏移量继续消费数据;查了资料发现kill -15 才能触发ShutdownHook(钩子:收到异常退出spark会自动清理相关连接对象,保证checkpoint完整;我还注册高优先级不过并没有什么用),这种方式不能新建StreamingContext不然每次都要从设置的偏移量(最早或最新)开始,这里我用的getOrCreate方法如下。---这种方式后面我发现我任务需要一直运行,目录下就会产生大量文件所以我就没用了,而且一直也没成功。

//使用 JavaStreamingContext.getOrCreate 自动处理
      JavaStreamingContext jssc = JavaStreamingContext.getOrCreate(checkpointLocation,
              () -> {
                  System.out.println("创建新的 StreamingContext");
                    JavaStreamingContext context = new JavaStreamingContext(conf, Durations.minutes(triggerInterval));
                   context.checkpoint(checkpointLocation);
                   return context;
               }
       );

3.这里不得不说又有个细节问题:SparkSession对象创建后自带一个SparkContext对象,所以StreamingContext不需要单独创建SparkContext,直接引用前面的就行,(这两个其实好像都是自带一个SparkContext,可以共用,重复创建会警告),这里我开始其实使用的SparkSession引用StreamingContext创建的SparkContext对象,但是报错找不到我的hive库,所以说这里大大的坑,SparkSession作为程序核心入口,它的SparkContext才具备hive相关的连接配置。

// 创建SparkSession
        SparkSession spark = SparkSession.builder()
                .config(conf)
                .enableHiveSupport()
                .getOrCreate();
        System.out.println("获取SparkSession成功");

        // 从SparkSession中获取已有的SparkContext
        // spark.sparkContext()返回的是Scala版本的SparkContext,转换为Java兼容的上下文
        JavaSparkContext javaSparkContext = JavaSparkContext.fromSparkContext(spark.sparkContext());
        JavaStreamingContext jssc = new JavaStreamingContext(javaSparkContext, Durations.seconds(triggerInterval));
        System.out.println("获取StreamingContext成功");

4.如果hive表有分区手动开启动态分区就好了

// 启用动态分区
        spark.conf().set("hive.exec.dynamic.partition", "true");
        spark.conf().set("hive.exec.dynamic.partition.mode", "nonstrict");

最后还遗留了一个问题 kafka没有生产数据,我程序一直挂着,设置60秒一个小批次,报了以下什么慢查询错误,也不知道啥原因:

Slow ReadProcessor read fields took 30081ms (threshold=30000ms); ack: seqno: 20 reply: SUCCESS reply: SUCCESS reply: SUCCESS downstreamAckTimeNanos: 20518849 flag: 0 flag: flag: 0, targets: [DatanodeInfoWithStorage[10.194.132.229:50010,DS-ba6c1cc2-3c86-49db-92c6-b8c69e16f8c7,DISK], DatanodeInfoWithStorage[10.194.134.20:50010,DS-349caef9-2d37-4850-b58b-4f05b4478183,DISK], DatanodeInfoWithStorage[10.194.132.250:50010,DS-239b0ac6-54f2-4f3d-beb2-ee0fe8340dc0,DISK]]

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐