【问题标题】:Amazon SQS Message AttributesAmazon SQS 消息属性
【发布时间】:2014-11-10 12:44:10
【问题描述】:

我无法正确添加属性。我正在使用 AWS 1.7 当我添加它们时,它们会显示在消息正文中,而不是属性中。当我登录 AWS 控制台时,我可以看到这一点。

我用这段代码添加消息属性:

                Message awsMessage = new Message();         

    Map<String, MessageAttributeValue> messageAttributes = 
new HashMap<String, MessageAttributeValue>();

                messageAttributes.put("email", new MessageAttributeValue()
                .withDataType("String")
                .withStringValue(email));
                messageAttributes.put("data", new MessageAttributeValue()
                .withDataType("String")
                .withStringValue(newFileName));
                messageAttributes.put("template", new MessageAttributeValue()
                .withDataType("String")
                .withStringValue(filename));

                awsMessage.setMessageAttributes(messageAttributes);

我尝试使用它来提取属性:

List<Message> messages = SQSUtilityClass.getMessagesFromQueue(QUEUE_URL);
    int size = messages.size();
    Map<String, String> attributes = new HashMap<String, String>();
    System.out.println("Size: "+size);
    for(int x =0;x<size;x++){
        Message message = messages.get(x);
        //System.out.println(message.getBody());
        attributes = message.getAttributes();
        for(String key: attributes.keySet()){
            System.out.println(key + " - "+ attributes.get(key));
        }
    }

当我通过 AWS 控制台查看时,我的属性仍然在消息正文中。

【问题讨论】:

    标签: java amazon-web-services message-queue amazon-sqs


    【解决方案1】:

    首先。您需要指定要随消息一起接收的属性。以下代码将接收属性emaildatatemplate(如您的示例所示):

    final AmazonSQS sqs = // your code to get SQS instance
    final String queue = // queue url
    final List<Message> messages = sqs.receiveMessage(
        new ReceiveMessageRequest(queue)
            .withMessageAttributeNames("email", "data", "template")
    ).getMessages();
    

    第二。你应该使用message.getMessageAttributes() 而不是message.getAttributes()

    final String email =
        message.getMessageAttributes().get("email").getStringValue()
    

    【讨论】:

      【解决方案2】:

      您可以使用以下库生成到 SQS。

        <dependency>
            <groupId>com.amazonaws</groupId>
            <artifactId>amazon-sqs-java-messaging-lib</artifactId>
            <version>${aws-messaging-lib-version}</version>
            <type>jar</type>
        </dependency>
      

      这是一个要生成的示例代码。

          SendMessageRequest request = new SendMessageRequest();
          request.withMessageAttributes(mapMessageAttributes(messageAttributes));
          request.setMessageBody(message);
          request.setQueueUrl(queueUrl);
          sqsClient.sendMessage(request);
      

      其中 mapMessageAttributes(messageAttributes) 返回一个 Map,而 sqsClient 是 AmazonSQSClient 的一个实例。

      一旦属性被传递到 SQS,我将使用 JMS 库通过将它们设置为 Spring DefaultMessageListenerContainer 来使用它们。

              SQSConnectionFactory connectionFactory =
                      SQSConnectionFactory.builder()
                              .withRegion(AWSSQSUtils.getAWSRegion(region))
                              .withAWSCredentialsProvider(new DefaultAWSCredentialsProviderChain())
                              .withNumberOfMessagesToPrefetch(prefetchCount)
                              .withClientConfiguration((new ClientConfiguration())
                                      .withMaxConnections(connectionsCount)
                                      .withProxyHost(proxyHost)
                                      .withProxyPort(proxyPort))
                              .build();
      
              LOG.info("Connection factory initialized..");
      
              // Create the connection.
              SQSConnection connection = connectionFactory.createConnection();
      
              // Create the transacted session with UNORDERED_ACKNOWLEDGE mode.
              Session session = connection.createSession(false, SQSSession.UNORDERED_ACKNOWLEDGE);
              Queue immediateQueue = session.createQueue(safInputQueueImmediate);
              Queue normalQueue = session.createQueue(safInputQueueNormal);
      
              LOG.info("Connection established and session created..");
      
              //Create listeners
              immediateListnerContainer = new
                      DefaultMessageListenerContainer();
              immediateListnerContainer.setConnectionFactory(connectionFactory);
              immediateListnerContainer.setMaxConcurrentConsumers(consumerConcurrency);
              immediateListnerContainer.setSessionAcknowledgeMode(SQSSession.UNORDERED_ACKNOWLEDGE);
              immediateListnerContainer.setCacheLevel(DefaultMessageListenerContainer.CACHE_CONSUMER);
              immediateListnerContainer.setDestination(immediateQueue);
              immediateListnerContainer.setMessageListener(
                      new StoreAndForwardListener(jobService, persistentMessagingDao));
              immediateListnerContainer.setMaxMessagesPerTask(maxMessagesPerTask);
              immediateListnerContainer.afterPropertiesSet();
              immediateListnerContainer.start();
      

      我使用的是程序化方法而不是声明性方法,因为 SQS 支持无序认知,而 Spring JMS 尚不支持。

      在MessageListener的onMessage方法中,你会得到一个javax.jms.Message对象,它可以用来获取所有的消息属性,如下所示。

      Enumeration<String> headers = message.getPropertyNames();
      while(headers.hasMoreElements()){
         attributeName = headers.nextElement();
         print(message.getStringProperty(attributeName))
      }
      

      希望对您有所帮助。

      我确实有一张使用 AWS 的票,其中一些属性在传输过程中被丢弃。即我设置了 M 个属性,只有 M-1 或 M-2 个属性才能进入 SQS。一有更新,我会在这里添加。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2014-09-26
        • 1970-01-01
        • 2014-06-27
        • 1970-01-01
        • 1970-01-01
        • 2016-03-05
        • 2015-09-09
        相关资源
        最近更新 更多