无法从本地Kafka Topic写入数据至GCS Bucket的技术求助
问题描述
在Dataproc集群上运行Spark Streaming作业,读取本地Kerberos认证的Kafka Topic数据,Schema Registry采用基础认证。日志显示已成功订阅Topic,但持续输出"无已提交偏移量、重置偏移量并查找最新偏移量"的内容,同时无法将数据写入GCS Bucket。
相关日志
Debug is true storeKey true useTicketCache false useKeyTab true doNotPrompt false ticketCache is null isInitiator true KeyTab is ./fid_name.keytab refreshKrb5Config is false principal is fid_name@domain.ca tryFirstPass is false useFirstPass is false storePass is false clearPass is false principal is fid_name@domain.ca Will use keytab Commit Succeeded
23/06/26 20:42:37 INFO org.apache.kafka.common.security.authenticator.AbstractLogin: Successfully logged in. 23/06/26 20:42:37 INFO org.apache.kafka.common.security.kerberos.KerberosLogin: [Principal=fid_name@domain.ca]: TGT refresh thread started. 23/06/26 20:42:37 INFO org.apache.kafka.common.security.kerberos.KerberosLogin: [Principal=fid_name@domain.ca]: TGT valid starting at: Mon Jun 26 20:42:37 UTC 2023 23/06/26 20:42:37 INFO org.apache.kafka.common.security.kerberos.KerberosLogin: [Principal=fid_name@domain.ca]: TGT expires: Tue Jun 27 06:42:37 UTC 2023 23/06/26 20:42:37 INFO org.apache.kafka.common.security.kerberos.KerberosLogin: [Principal=fid_name@domain.ca]: TGT refresh sleeping until: Tue Jun 27 04:58:31 UTC 2023 23/06/26 20:42:38 INFO org.apache.kafka.common.utils.AppInfoParser: Kafka version: 2.6.0 23/06/26 20:42:38 INFO org.apache.kafka.common.utils.AppInfoParser: Kafka commitId: 62ab7664f6e03875 23/06/26 20:42:38 INFO org.apache.kafka.common.utils.AppInfoParser: Kafka startTimeMs: 1687812158067 23/06/26 20:42:38 INFO org.apache.kafka.clients.consumer.KafkaConsumer: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Subscribed to topic(s): topic_1 23/06/26 20:42:38 INFO org.apache.kafka.clients.Metadata: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Cluster ID: V5F9mOvYTxSZdi3O6QF4lQ 23/06/26 20:42:38 INFO org.apache.kafka.clients.consumer.internals.AbstractCoordinator: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Discovered group coordinator broker_1:9093 (id: 2147483641 rack: null) 23/06/26 20:42:38 INFO org.apache.kafka.clients.consumer.internals.AbstractCoordinator: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] (Re-)joining group 23/06/26 20:42:38 INFO org.apache.kafka.clients.consumer.internals.AbstractCoordinator: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Join group failed with org.apache.kafka.common.errors.MemberIdRequiredException: The group member needs to have a valid member id before actually entering a consumer group. 23/06/26 20:42:38 INFO org.apache.kafka.clients.consumer.internals.AbstractCoordinator: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] (Re-)joining group 23/06/26 20:42:41 INFO org.apache.kafka.clients.consumer.internals.ConsumerCoordinator: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Finished assignment for group at generation 1: {consumer-consumger_group_name-1-f1070ddf-1ff0-4857-9664-50483f233a96=Assignment(partitions=[topic_1-0, topic_1-1, topic_1-2, topic_1-3, topic_1-4, topic_1-5])} 23/06/26 20:42:41 INFO org.apache.kafka.clients.consumer.internals.AbstractCoordinator: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Successfully joined group with generation 1 23/06/26 20:42:41 INFO org.apache.kafka.clients.consumer.internals.ConsumerCoordinator: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Notifying assignor about the new Assignment(partitions=[topic_1-0, topic_1-1, topic_1-2, topic_1-3, topic_1-4, topic_1-5]) 23/06/26 20:42:41 INFO org.apache.kafka.clients.consumer.internals.ConsumerCoordinator: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Adding newly assigned partitions: topic_1-2, topic_1-1, topic_1-4, topic_1-3, topic_1-0, topic_1-5 23/06/26 20:42:41 INFO org.apache.kafka.clients.consumer.internals.ConsumerCoordinator: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Found no committed offset for partition topic_1-2 23/06/26 20:42:41 INFO org.apache.kafka.clients.consumer.internals.ConsumerCoordinator: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Found no committed offset for partition topic_1-1 23/06/26 20:42:41 INFO org.apache.kafka.clients.consumer.internals.ConsumerCoordinator: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Found no committed offset for partition topic_1-4 23/06/26 20:42:41 INFO org.apache.kafka.clients.consumer.internals.ConsumerCoordinator: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Found no committed offset for partition topic_1-3 23/06/26 20:42:41 INFO org.apache.kafka.clients.consumer.internals.ConsumerCoordinator: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Found no committed offset for partition topic_1-0 23/06/26 20:42:41 INFO org.apache.kafka.clients.consumer.internals.ConsumerCoordinator: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Found no committed offset for partition topic_1-5 23/06/26 20:42:42 INFO org.apache.kafka.clients.consumer.internals.SubscriptionState: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Resetting offset for partition topic_1-3 to offset 0. 23/06/26 20:42:42 INFO org.apache.kafka.clients.consumer.internals.SubscriptionState: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Resetting offset for partition topic_1-2 to offset 0. 23/06/26 20:42:42 INFO org.apache.kafka.clients.consumer.internals.SubscriptionState: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Resetting offset for partition topic_1-0 to offset 0. 23/06/26 20:42:42 INFO org.apache.kafka.clients.consumer.internals.SubscriptionState: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Resetting offset for partition topic_1-1 to offset 0. 23/06/26 20:42:42 INFO org.apache.kafka.clients.consumer.internals.SubscriptionState: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Resetting offset for partition topic_1-5 to offset 0. 23/06/26 20:42:42 INFO org.apache.kafka.clients.consumer.internals.SubscriptionState: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Resetting offset for partition topic_1-4 to offset 1. 23/06/26 20:42:42 INFO org.apache.kafka.clients.consumer.internals.SubscriptionState: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Seeking to LATEST offset of partition topic_1-2 23/06/26 20:42:42 INFO org.apache.kafka.clients.consumer.internals.SubscriptionState: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Seeking to LATEST offset of partition topic_1-1 23/06/26 20:42:42 INFO org.apache.kafka.clients.consumer.internals.SubscriptionState: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Seeking to LATEST offset of partition topic_1-4 23/06/26 20:42:42 INFO org.apache.kafka.clients.consumer.internals.SubscriptionState: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Seeking to LATEST offset of partition topic_1-3 23/06/26 20:42:42 INFO org.apache.kafka.clients.consumer.internals.SubscriptionState: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Seeking to LATEST offset of partition topic_1-0 23/06/26 20:42:42 INFO org.apache.kafka.clients.consumer.internals.SubscriptionState: [Consumer clientId=consumer-consumger_group_name-1, groupId=consumger_group_name] Seeking to LATEST offset of partition topic_1-5
作业提交命令
gcloud dataproc jobs submit spark --project=prj-cxbi-dev --cluster=multi-node-cluster --region=northamerica-northeast1 --class=Logr.AvroConsumer --jars=xyz_project.jar --properties="spark.submit.deployMode"="client","spark.dynamicAllocation.enabled"="true","spark.shuffle.service.enabled"="true","spark.executor.memory"="12g","spark.driver.memory"="4g","spark.executor.cores"="5","spark.driver.extraJavaOptions=-Djava.security.auth.login.config=kafka_np_client_jaas.conf","spark.executor.extraJavaOptions=-Djava.security.auth.login.config=kafka_np_client_jaas.conf","spark.driver.extraJavaOptions=-Djavax.net.ssl.trustStore=truststore.jks","spark.executor.extraJavaOptions=-Djavax.net.ssl.trustStore=truststore.jks" --files kafka_np_client_jaas.conf,truststore.jks,fid_np.keytab -- "topic_name" "schema_version" "90 seconds" "gs://checkpoint" "gs://bucket_where_data_needs_to_be_written" "job_name"
排查与解决方法
一、Kafka偏移量问题处理
验证消费者组偏移量状态
- 使用Kafka命令行工具查询指定消费者组的偏移量提交情况:
kafka-consumer-groups.sh --bootstrap-server <kafka_broker>:9093 --command-config <kerberos_config_file> --group consumger_group_name --describe - 若确实无提交记录,检查作业的偏移量提交配置:
- Spark Streaming Direct Stream需确认
enable.auto.commit是否设为true; - 若使用手动提交,需确保代码中调用了
commitAsync()或commitSync()方法。
- Spark Streaming Direct Stream需确认
- 使用Kafka命令行工具查询指定消费者组的偏移量提交情况:
调整偏移量重置策略
- 日志显示重置到最新偏移量,若需要从头消费,可将
auto.offset.reset配置改为earliest; - 若重复运行仍无偏移量提交,重点排查代码中的偏移量提交逻辑是否存在遗漏或异常。
- 日志显示重置到最新偏移量,若需要从头消费,可将
修复消费者组加入异常
- 日志中出现
MemberIdRequiredException,可能是版本兼容或Kerberos会话问题:- 确认Spark与Kafka版本兼容性(Spark 3.x建议搭配Kafka 2.4+);
- 检查Kerberos TGT有效期(当前为10小时),确保作业运行时长不超过有效期,或配置自动刷新票据。
- 日志中出现
二、GCS写入失败排查
检查集群GCS权限
- 确认Dataproc集群的服务账号是否拥有目标GCS Bucket的写入权限(
storage.objects.create、storage.objects.list等),可通过IAM控制台查看权限配置。
- 确认Dataproc集群的服务账号是否拥有目标GCS Bucket的写入权限(
验证GCS路径配置
- 确认命令中传入的
gs://bucket_where_data_needs_to_be_written路径正确,Bucket存在且无拼写错误。
- 确认命令中传入的
检查Spark-GCS集成配置
- Dataproc默认集成GCS连接器,若作业覆盖了相关配置,需添加:
spark.hadoop.fs.gs.impl=com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem - 排查SSL配置是否干扰GCS连接,可临时关闭SSL验证测试。
- Dataproc默认集成GCS连接器,若作业覆盖了相关配置,需添加:
提取GCS写入日志
- 查看YARN或Dataproc集群的Executor日志,查找GCS写入相关的异常(如权限拒绝、连接超时)。
三、修正作业配置问题
合并重复的JVM配置
- 原命令中
spark.driver.extraJavaOptions和spark.executor.extraJavaOptions被重复设置,导致前序配置被覆盖,需合并为:--properties=...,spark.driver.extraJavaOptions="-Djava.security.auth.login.config=kafka_np_client_jaas.conf -Djavax.net.ssl.trustStore=truststore.jks",spark.executor.extraJavaOptions="-Djava.security.auth.login.config=kafka_np_client_jaas.conf -Djavax.net.ssl.trustStore=truststore.jks"
- 原命令中
验证JAAS配置文件
- 确保
kafka_np_client_jaas.conf中的principal和keytab路径与传入的文件一致,示例配置:KafkaClient { com.sun.security.auth.module.Krb5LoginModule required useKeyTab=true keyTab="./fid_np.keytab" principal="fid_name@domain.ca" storeKey=true useTicketCache=false; };
- 确保
内容的提问来源于stack exchange,提问作者S D
相关产品推荐
相关产品推荐

