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

Spark Structured Streaming基于SASL OAuthBearer连接Spring Kafka的实现问询

Spark Structured Streaming 连接 Kafka 用 SASL_SSL/OAuthBearer 实现(Java/Scala)

Spark Structured Streaming 底层复用 Kafka Java 客户端进行通信,所以只需把 OAuthBearer 相关配置传递给 Spark 的 Kafka 数据源,同时实现自定义逻辑从 ADSF 门户获取令牌即可。下面是具体实现方案:

核心配置参数

首先要在 Spark 的 Kafka 配置里指定 SASL 相关项:

  • security.protocol: 设为 SASL_SSL
  • sasl.mechanism: 设为 OAUTHBEARER
  • sasl.jaas.config: 指定 OAuthBearer 的 JAAS 配置,关联自定义的令牌回调处理器
  • 额外 SSL 配置(如信任库路径、密码)根据 Kafka 集群要求补充

Java 版本实现

1. 自定义 OAuthBearer 回调处理器

实现 CallbackHandler,在里面调用 ADSF 接口拿令牌:

import org.apache.kafka.common.security.oauthbearer.OAuthBearerTokenCallback;
import javax.security.auth.callback.Callback;
import javax.security.auth.callback.UnsupportedCallbackException;
import javax.security.auth.login.AppConfigurationEntry;
import java.io.IOException;
import java.util.Map;

public class ADSFOAuthCallbackHandler implements javax.security.auth.callback.CallbackHandler {

    @Override
    public void handle(Callback[] callbacks) throws IOException, UnsupportedCallbackException {
        for (Callback callback : callbacks) {
            if (callback instanceof OAuthBearerTokenCallback) {
                OAuthBearerTokenCallback oauthCallback = (OAuthBearerTokenCallback) callback;
                // 调用内部 ADSF 门户接口获取令牌,替换成你的实际逻辑
                String token = fetchADSFToken();
                oauthCallback.token(token);
            } else {
                throw new UnsupportedCallbackException(callback, "不支持的回调类型");
            }
        }
    }

    private String fetchADSFToken() throws IOException {
        // 示例:这里写调用 ADSF 令牌接口的代码,比如用 HttpClient 发请求
        // 记得处理令牌过期刷新逻辑
        return "ADSF返回的有效令牌";
    }

    // 可选:处理 JAAS 配置中的自定义参数
    public void configure(Map<String, String> configs, String mechanism, AppConfigurationEntry[] jaasConfigEntries) {
        // 比如从configs里读取ADSF的端点地址等参数
    }
}

2. Spark 流式程序代码

把 Kafka 配置传入 Spark 数据源,指定回调处理器:

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.streaming.StreamingQueryException;
import java.util.HashMap;
import java.util.Map;

public class SparkKafkaOAuthDemo {
    public static void main(String[] args) throws StreamingQueryException {
        SparkSession spark = SparkSession.builder()
                .appName("Spark-Kafka-OAuth-Demo")
                .master("local[*]") // 生产环境删除该行
                .getOrCreate();

        Map<String, String> kafkaConfigs = new HashMap<>();
        kafkaConfigs.put("kafka.bootstrap.servers", "你的Kafka broker地址:9093");
        kafkaConfigs.put("subscribe", "目标主题名");
        kafkaConfigs.put("security.protocol", "SASL_SSL");
        kafkaConfigs.put("sasl.mechanism", "OAUTHBEARER");
        // JAAS配置:指定自定义回调处理器的全类名
        kafkaConfigs.put("sasl.jaas.config", "org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required " +
                "callback.handler.class=\"com.xxx.ADSFOAuthCallbackHandler\";");
        // SSL信任库配置(按需添加)
        kafkaConfigs.put("ssl.truststore.location", "/path/to/truststore.jks");
        kafkaConfigs.put("ssl.truststore.password", "信任库密码");

        // 读取Kafka流
        Dataset<Row> kafkaStream = spark.readStream()
                .format("kafka")
                .options(kafkaConfigs)
                .load();

        // 简单处理并输出到控制台
        kafkaStream.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
                .writeStream()
                .outputMode("append")
                .format("console")
                .start()
                .awaitTermination();
    }
}

Scala 版本实现

1. 自定义回调处理器

import org.apache.kafka.common.security.oauthbearer.OAuthBearerTokenCallback
import javax.security.auth.callback.{Callback, UnsupportedCallbackException}
import javax.security.auth.login.AppConfigurationEntry
import java.io.IOException
import java.util.Map

class ADSFOAuthCallbackHandler extends javax.security.auth.callback.CallbackHandler {

  override def handle(callbacks: Array[Callback]): Unit = {
    callbacks.foreach {
      case oauthCallback: OAuthBearerTokenCallback =>
        val token = fetchADSFToken()
        oauthCallback.token(token)
      case callback =>
        throw new UnsupportedCallbackException(callback, "不支持的回调类型")
    }
  }

  @throws[IOException]
  private def fetchADSFToken(): String = {
    // 实现调用ADSF获取令牌的逻辑,注意处理过期刷新
    "ADSF返回的有效令牌"
  }

  override def configure(configs: Map[String, String], mechanism: String, jaasConfigEntries: Array[AppConfigurationEntry]): Unit = {
    // 可选:处理自定义配置参数
  }
}

2. Spark 流式程序代码

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.streaming.StreamingQueryException

object SparkKafkaOAuthDemo {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder()
      .appName("Spark-Kafka-OAuth-Demo")
      .master("local[*]") // 生产环境移除
      .getOrCreate()

    import spark.implicits._

    val kafkaConfigs = Map(
      "kafka.bootstrap.servers" -> "你的Kafka broker地址:9093",
      "subscribe" -> "目标主题名",
      "security.protocol" -> "SASL_SSL",
      "sasl.mechanism" -> "OAUTHBEARER",
      "sasl.jaas.config" -> "org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required " +
        "callback.handler.class=\"com.xxx.ADSFOAuthCallbackHandler\";",
      "ssl.truststore.location" -> "/path/to/truststore.jks",
      "ssl.truststore.password" -> "信任库密码"
    )

    val kafkaStream = spark.readStream
      .format("kafka")
      .options(kafkaConfigs)
      .load()

    kafkaStream.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
      .writeStream
      .outputMode("append")
      .format("console")
      .start()
      .awaitTermination()
  }
}

注意事项

  • 打包时要把自定义回调处理器和依赖的 HTTP 客户端等包打进 Spark 应用 jar 里,确保类路径能加载到。
  • ADSF 令牌有过期时间,回调处理器里要实现自动刷新逻辑,避免令牌过期导致连接中断。
  • 生产环境里敏感信息(如信任库密码、ADSF 认证参数)不要硬编码,用环境变量或密钥管理服务传递。
  • 可以先用 kafka-console-consumer.sh 测试 OAuthBearer 配置是否能正常消费,再迁移到 Spark 程序。

内容的提问来源于stack exchange,提问作者novice programmer

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 00:28:27