【问题标题】:Why spring-cloud-stream not populating JMSMessageID while publishing messages to Solace topics?为什么 spring-cloud-stream 在向 Solace 主题发布消息时不填充 JMSMessageID?
【发布时间】:2019-09-10 01:59:24
【问题描述】:

问题总结

在我的项目中,我们正在尝试使用 Spring-Cloud-Stream (SCS) 连接到 Solace。最终我们计划搬到卡夫卡。因此,使用 SCS 将帮助我们非常轻松地迁移到 Kafka,而无需任何代码更改以及极少的配置和依赖项更改。

我们已经使用 JMS 使用 Solace 有一段时间了。现在,当我们尝试使用 SCS 向 Solace 发布消息时,我们观察到在消息中,一些关键的 JMS 标头(JMSMessageID、JMSType、JMSPriority、JMSCorrelationID、JMSExpiration)是空白的。

我们需要单独配置 JMS 标头吗?如果是,如何?

我已经尝试过的

我试图设置这样的标题,但这只会导致具有相同名称的重复标题。

    @Output(SendReport.TO_NMR)
    public void sendMessage(String request) {
        log.info("****************** Got this Report Request: " + request);

        MessageBuilder<String> builder = MessageBuilder.withPayload(request);
        builder.setHeader("JMSType","report-request");
        builder.setHeader("JMSMessageId","1");
        builder.setHeader("JMSCorrelationId","11");
        builder.setHeader("JMSMessageID","4");
        builder.setHeader("JMSCorrelationID","114");
        builder.setHeader("ApplicationMessageId","111");
        builder.setHeader("ApplicationMessageID","112");
        builder.setCorrelationId("23434");

        Message message = builder.build();
        sendReport.output().send(message);
    }

Solace 中消息的 JMS 标头如下所示

JMSMessageID    
JMSDestination  TOPIC_NAME
JMSTimestamp    Wed Dec 31 18:00:00 CST 1969
JMSType 
JMSReplyTo  
JMSCorrelationID    
JMSExpiration   0
JMSPriority 0
JMSType nmr-report-request
JMSMessageId    1
JMSMessageID    4
_isJavaSerializedObject-contentType true
_isJavaSerializedObject-id  true
solaceSpringCloudStreamBinderVersion    0.1.0
ApplicationMessageId    111
ApplicationMessageID    112
JMSCorrelationId    11
JMSCorrelationID    114
correlationId   23434
id  [-84,-19,0,5,115,114,0,14,106,97,118,97,46,117,116,105,108,46,85,85,73,68,-68,-103,3,-9,-104,109,-123,47,2,0,2,74,0,12,108,101,97,115,116,83,105,103,66,105,116,115,74,0,11,109,111,115,116,83,105,103,66,105,116,115,120,112,13,-26,2,-51,111,-17,73,73,-18,-32,-26,-11,-46,-89,50,-37] (offset=377, length=80)
contentType [-84,-19,0,5,115,114,0,33,111,114,103,46,115,112,114,105,110,103,102,114,97,109,101,119,111,114,107,46,117,116,105,108,46,77,105,109,101,84,121,112,101,56,-76,29,-63,64,96,-36,-81,2,0,3,76,0,10,112,97,114,97,109,101,116,101,114,115,116,0,15,76,106,97,118,97,47,117,116,105,108,47,77,97,112,59,76,0,7,115,117,98,116,121,112,101,116,0,18,76,106,97,118] (offset=473, length=190)
timestamp   1555707627482

用于连接 Solace 的代码

Spring Boot 主类

@SpringBootApplication
@EnableDiscoveryClient
@Slf4j
@EnableBinding({SendReport.class}) 
public class ReportServerApplication {

    public static void main(final String[] args) {
        ApplicationContext ctx = new ClassPathXmlApplicationContext("applicationContext-server.xml");
        new SpringApplicationBuilder(ReportServerApplication.class).listeners(new EnvironmentPreparedListener())                                                   .run(args);
}

将频道连接到主题的类:

public interface SendReport {

    String TO_NMR = "solace-poc-outbound";

    @Output(SendReport.TO_NMR)
    MessageChannel output();

}

消息处理程序:

@Slf4j
@Component
@EnableBinding({SendReport.class})
public class MessageHandler {

    private SendReport sendReport;

    public MessageHandler(SendReport sendReport){
        this.sendReport = sendReport;
    }

    @Output(SendReport.TO_NMR)
    public void sendMessage(String request) {
        log.info("****************** Got this Report Request: " + request);
        var message = MessageBuilder.withPayload(request).build();
        sendReport.output().send(message);
    }
}

用于配置的属性:application.yml

spring:
  cloud:
    # spring cloud stream binding
    stream:
      bindings:
        solace-poc-outbound:
          destination: TOPIC_NAME
          contentType: text/plain

solace:
  java:
    host: tcp://xyz.abc.com
    #port: xxx
    msgVpn: yyy
    clientUsername: aaa

使用的依赖项:

'org.springframework.cloud:spring-cloud-stream',
'com.solace.spring.cloud:spring-cloud-starter-stream-solace:1.1.+'

观察

  • 预期结果:所有 JMS 标头都应由 SCS 填充。
  • 实际结果:某些 JMS 标头未填充。

【问题讨论】:

    标签: jms spring-cloud-stream solace


    【解决方案1】:

    参见 JMS MessageJavaDocs:

    /** Sets the message ID.
      *  
     * <P>This method is for use by JMS providers only to set this field 
     * when a message is sent. This message cannot be used by clients 
     * to configure the message ID. This method is public
     * to allow a JMS provider to set this field when sending a message
     * whose implementation is not its own.
      *
      * @param id the ID of the message
      *
      * @exception JMSException if the JMS provider fails to set the message ID 
      *                         due to some internal error.
      *
      * @see javax.jms.Message#getJMSMessageID()
      */ 
    
    void
    setJMSMessageID(String id) throws JMSException;
    

    因此,无法从应用程序级别填充此属性。

    在 ActiveMQ 中我看到这样的代码:

     msg.setMessageId(new MessageId(producer.getProducerInfo().getProducerId(), sequenceNumber));
    
     // Set the message id.
     if (msg != message) {
            message.setJMSMessageID(msg.getMessageId().toString());
    

    但仍然:这不是我们可以从应用程序级别控制的。

    prioritydeliveryModetimeToLive 可以从 JmsSendingMessageHandler 填充:

    if (this.jmsTemplate instanceof DynamicJmsTemplate && this.jmsTemplate.isExplicitQosEnabled()) {
            Integer priority = StaticMessageHeaderAccessor.getPriority(message);
            if (priority != null) {
                DynamicJmsTemplateProperties.setPriority(priority);
            }
            if (this.deliveryModeExpression != null) {
                Integer deliveryMode =
                        this.deliveryModeExpression.getValue(this.evaluationContext, message, Integer.class);
    
                if (deliveryMode != null) {
                    DynamicJmsTemplateProperties.setDeliveryMode(deliveryMode);
                }
            }
            if (this.timeToLiveExpression != null) {
                Long timeToLive = this.timeToLiveExpression.getValue(this.evaluationContext, message, Long.class);
                if (timeToLive != null) {
                    DynamicJmsTemplateProperties.setTimeToLive(timeToLive);
                }
            }
        }
    

    JmsCorrelationID 必须由 JmsHeaders.CORRELATION_ID 填充。 JmsType分别由JmsHeaders.TYPE

    public void fromHeaders(MessageHeaders headers, javax.jms.Message jmsMessage) {
        try {
            Object jmsCorrelationId = headers.get(JmsHeaders.CORRELATION_ID);
            if (jmsCorrelationId instanceof Number) {
                jmsCorrelationId = jmsCorrelationId.toString();
            }
            if (jmsCorrelationId instanceof String) {
                try {
                    jmsMessage.setJMSCorrelationID((String) jmsCorrelationId);
                }
                catch (Exception e) {
                    this.logger.info("failed to set JMSCorrelationID, skipping", e);
                }
            }
            Object jmsReplyTo = headers.get(JmsHeaders.REPLY_TO);
            if (jmsReplyTo instanceof Destination) {
                try {
                    jmsMessage.setJMSReplyTo((Destination) jmsReplyTo);
                }
                catch (Exception e) {
                    this.logger.info("failed to set JMSReplyTo, skipping", e);
                }
            }
            Object jmsType = headers.get(JmsHeaders.TYPE);
            if (jmsType instanceof String) {
                try {
                    jmsMessage.setJMSType((String) jmsType);
                }
                catch (Exception e) {
                    this.logger.info("failed to set JMSType, skipping", e);
                }
            }
    

    更多信息请参见DefaultJmsHeaderMapper

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2020-12-26
      • 2016-09-09
      • 2018-07-18
      • 2021-06-12
      • 1970-01-01
      • 2022-01-06
      • 2020-09-18
      相关资源
      最近更新 更多