You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

无法从本地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偏移量问题处理

  1. 验证消费者组偏移量状态

    • 使用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()方法。
  2. 调整偏移量重置策略

    • 日志显示重置到最新偏移量,若需要从头消费,可将auto.offset.reset配置改为earliest;
    • 若重复运行仍无偏移量提交,重点排查代码中的偏移量提交逻辑是否存在遗漏或异常。
  3. 修复消费者组加入异常

    • 日志中出现MemberIdRequiredException,可能是版本兼容或Kerberos会话问题:
      • 确认Spark与Kafka版本兼容性(Spark 3.x建议搭配Kafka 2.4+);
      • 检查Kerberos TGT有效期(当前为10小时),确保作业运行时长不超过有效期,或配置自动刷新票据。

二、GCS写入失败排查

  1. 检查集群GCS权限

    • 确认Dataproc集群的服务账号是否拥有目标GCS Bucket的写入权限(storage.objects.create、storage.objects.list等),可通过IAM控制台查看权限配置。
  2. 验证GCS路径配置

    • 确认命令中传入的gs://bucket_where_data_needs_to_be_written路径正确,Bucket存在且无拼写错误。
  3. 检查Spark-GCS集成配置

    • Dataproc默认集成GCS连接器,若作业覆盖了相关配置,需添加:
      spark.hadoop.fs.gs.impl=com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem
      
    • 排查SSL配置是否干扰GCS连接,可临时关闭SSL验证测试。
  4. 提取GCS写入日志

    • 查看YARN或Dataproc集群的Executor日志,查找GCS写入相关的异常(如权限拒绝、连接超时)。

三、修正作业配置问题

  1. 合并重复的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"
      
  2. 验证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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.17 21:26:59