TLS/SSL 在 Flink/Kafka 数据传输中的配置实现

在大数据生态中,Flink 和 Kafka 的数据传输安全至关重要。通过 TLS/SSL 加密可确保数据传输的机密性完整性。以下是具体配置步骤:

一、核心概念
  1. TLS/SSL 作用

    • 加密数据传输通道
    • 身份验证(单向/双向认证)
    • 防止中间人攻击
  2. 关键组件

    • 密钥库(Keystore):存储服务端私钥和证书
    • 信任库(Truststore):存储可信 CA 证书
    • 证书格式:推荐 JKS 或 PKCS12

二、Kafka Broker 配置

步骤:

  1. 生成证书

    # 生成CA根证书
    openssl req -new -x509 -keyout ca-key -out ca-cert -days 365
    
    # 生成Broker密钥对
    keytool -keystore kafka.server.keystore.jks -alias broker -validity 365 -genkey
    

  2. 签名证书

    keytool -keystore kafka.server.keystore.jks -alias broker -certreq -file cert-req
    openssl x509 -req -CA ca-cert -CAkey ca-key -in cert-req -out cert-signed -days 365 -CAcreateserial
    keytool -keystore kafka.server.keystore.jks -alias broker -import -file cert-signed
    

  3. 配置 server.properties

    listeners=SSL://:9093
    security.inter.broker.protocol=SSL
    ssl.keystore.location=/path/to/kafka.server.keystore.jks
    ssl.keystore.password=your_password
    ssl.key.password=your_password
    ssl.truststore.location=/path/to/kafka.server.truststore.jks
    ssl.truststore.password=your_password
    ssl.client.auth=required  # 启用双向认证
    


三、Flink 连接 Kafka 配置

场景: Flink 作为 Kafka 消费者/生产者
Java 代码示例:

Properties props = new Properties();
props.setProperty("bootstrap.servers", "kafka-host:9093");
props.setProperty("security.protocol", "SSL");
props.setProperty("ssl.truststore.location", "/path/to/client.truststore.jks");
props.setProperty("ssl.truststore.password", "client_pass");
props.setProperty("ssl.keystore.location", "/path/to/client.keystore.jks"); // 双向认证需配置
props.setProperty("ssl.keystore.password", "client_pass");

// 创建消费者
FlinkKafkaConsumer<String> consumer = new FlinkKafkaConsumer<>(
    "topic", 
    new SimpleStringSchema(), 
    props
);

// 创建生产者
FlinkKafkaProducer<String> producer = new FlinkKafkaProducer<>(
    "topic",
    new SimpleStringSchema(),
    props
);


四、双向认证配置(可选)
  1. 客户端生成证书

    keytool -keystore client.keystore.jks -alias client -validity 365 -genkey
    # 签名步骤同Broker
    

  2. Kafka 配置
    确保 ssl.client.auth=required

  3. Flink 添加配置

    props.setProperty("ssl.keystore.location", "/path/to/client.keystore.jks");
    props.setProperty("ssl.keystore.password", "client_pass");
    


五、验证与调试
  1. 测试命令

    openssl s_client -connect kafka-host:9093 -showcerts
    

  2. 常见问题

    • 证书路径权限问题
    • 密码不一致
    • 证书过期(可通过 keytool -list -v -keystore xxx.jks 检查)

六、最佳实践
  1. 证书管理

    • 使用正式 CA 签发证书(非自签名)
    • 定期轮转证书(推荐 90 天)
  2. 性能优化

    • 启用 OpenSSL 原生加速(配置 ssl.provider=OPENSSL
    • 使用会话复用减少握手开销
  3. 安全增强

    ssl.enabled.protocols=TLSv1.3,TLSv1.2  # 禁用旧协议
    ssl.cipher.suites=TLS_AES_256_GCM_SHA384  # 强加密套件
    

:完整配置需保持 Kafka 集群、Flink 作业、所有客户端(如 Schema Registry)的 TLS 配置一致,避免中断。生产环境建议通过 ConfigMap(K8s)或 Ansible 统一管理配置。

更多推荐