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_SSLsasl.mechanism: 设为OAUTHBEARERsasl.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
相关产品推荐
相关产品推荐

