如何在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
相关产品推荐
相关产品推荐

