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

如何在Apache Beam应用中限制Qpid JMS连接Solace的数量?

如何减少Apache Beam应用中JMS连接的创建数量?

我的Apache Beam应用通过Apache Qpid JMS库及Java消息服务(JMS)向Solace队列发送消息。启动应用后,两个节点创建了2000多个连接,日志中频繁出现如下记录:

Connection ID:0030887-f9c8-429f-b5d1-eda15330e262:1 connected to server: amqps://my.queue.address:54321

但Solace的连接数限制为400个,这给我带来了困扰。


我的代码

调用写入器的代码

@NotNull
@Override
public POutput expand(@NotNull PCollection<String> input) {
  input.apply(
      "WriteMessageToMyTopic",
      JmsIOWriter.writeMessage());

  return PDone.in(input.getPipeline());
}

JMS Writer实现

package com.example;

import com.example.StaticConnectionFactory;
import java.io.Serializable;
import javax.jms.JMSException;
import javax.jms.Message;
import javax.jms.Session;
import org.apache.beam.sdk.io.jms.JmsIO;
import org.apache.beam.sdk.io.jms.RetryConfiguration;

final class JmsIOWriter implements Serializable {

  public static JmsIO.Write<String> writeMessage() {
    var connectionFactory = StaticConnectionFactory.getFactory();

    return JmsIO.<String>write()
        .withRetryConfiguration(RetryConfiguration.create(5))
        .withConnectionFactory(connectionFactory)
        .withTopicNameMapper(m -> topicMapper(m, "my/topic/"))
        .withValueMapper(JmsIOWriter::valueMapper);
  }

  private static Message valueMapper(String element, Session session) {
    try {
      return session.createTextMessage(element);
    } catch (JMSException ex) {
      throw new IllegalStateException(
          String.format("Unable to send message in queue %s", element), ex);
    }
  }

  private static String topicMapper(String message) {
    return "my/topic/" + message.length;
  }
}

静态连接工厂

package com.example;

import org.apache.qpid.jms.JmsConnectionFactory;

public final class StaticConnectionFactory {

  private static JmsConnectionFactory factory;

  public static synchronized JmsConnectionFactory getFactory() {
    if (factory == null) {
      factory = new SslJmsConnectionFactory();
      factory.setUsername("myUsername");
      factory.setPassword("myPassword");
      factory.setRemoteURI("amqps://my.queue.address:54321");
      factory.setForceAsyncAcks(true);
      factory.setReceiveLocalOnly(true);
      factory.setReceiveNoWaitLocalOnly(true);
    }
    return factory;
  }
}

SSL连接工厂

package com.example;

import static com.example.SslContextFactory.createSslContext;

import java.util.Base64;
import javax.jms.Connection;
import javax.jms.JMSException;

public class SslJmsConnectionFactory extends org.apache.qpid.jms.JmsConnectionFactory {

  @Override
  public Connection createConnection(String username, String password) throws JMSException {
    initSslContext();
    return super.createConnection(username, password);
  }

  private void initSslContext() {
    setSslContext(
        createSslContext(
            Base64.getDecoder().decode(this.getUsername()), this.getPassword().toCharArray()));
  }
}

SSL上下文工厂

package com.example;

import java.io.ByteArrayInputStream;
import java.io.IOException;
import java.security.KeyManagementException;
import java.security.KeyStore;
import java.security.KeyStoreException;
import java.security.NoSuchAlgorithmException;
import java.security.SecureRandom;
import java.security.UnrecoverableKeyException;
import java.security.cert.CertificateException;
import javax.net.ssl.KeyManagerFactory;
import javax.net.ssl.SSLContext;

public final class SslContextFactory {

  public static SSLContext createSslContext(byte[] credentials, char[] password) {

    try (var inputStream = new ByteArrayInputStream(credentials)) {
      var keyStore = KeyStore.getInstance("PKCS12");
      keyStore.load(inputStream, password);
      return initSslContext(keyStore, password);
    } catch (IOException | KeyStoreException | NoSuchAlgorithmException | CertificateException e) {
      throw new IllegalStateException("Cannot create SSL context", e);
    }
  }

  private static SSLContext initSslContext(KeyStore keyStore, char[] password) {
    try {
      var sslContext = SSLContext.getInstance("TLSv1.2");
      var keyManagerFactory = createKeyManagerFactory(keyStore, password);
      sslContext.init(keyManagerFactory.getKeyManagers(), null, new SecureRandom());
      return sslContext;
    } catch (NoSuchAlgorithmException | KeyManagementException e) {
      throw new IllegalStateException("Unable to initialize SSL context", e);
    }
  }

  private static KeyManagerFactory createKeyManagerFactory(KeyStore keyStore, char[] password) {
    try {
      var keyManagerFactory = KeyManagerFactory.getInstance("SunX509");
      keyManagerFactory.init(keyStore, password);
      return keyManagerFactory;
    } catch (NoSuchAlgorithmException | KeyStoreException | UnrecoverableKeyException e) {
      throw new IllegalStateException("Unable to initialize a key manager factory", e);
    }
  }
}

已尝试的方案

  • 使用静态JMS连接工厂避免重复创建连接,但无效果
  • 尝试在初始化JMS连接工厂时设置SSL上下文而非创建连接时设置,却出现认证错误:org.apache.qpid.jms.provider.exceptions.ProviderConnectionSecuritySaslException: Client failed to authenticate using SASL: PLAIN
  • 暂无法使用Apache Beam 2.61.0及以上版本提供的实验性Solace Writer

内容的提问来源于stack exchange,提问作者Vincent Doba

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 16:13:11