使用Spark Streaming 消费kafka数据到hive问题记录
记录首次使用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]]
更多推荐


所有评论(0)