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

无法从AWS SQS死信队列消费消息但主队列正常的排查求助

问题描述

使用Kotlin编写了两个独立的消息监听器类,分别监听两个AWS SQS标准队列(非FIFO):一个为主队列,另一个为死信队列(DLQ)。目前遇到的问题是:DLQ的消息监听器无法从对应DLQ中消费消息,但将DLQ名称替换为主队列名称时,该监听器可以正常消费;且AWS控制台显示DLQ中确实存在未消费的消息。怀疑是DLQ配置存在问题但未发现异常,需要根因排查的参考或提示。

初始化SQS连接代码

@PostConstruct
fun init() {
    try {
        if (enableCalmQueue) {
            initializeListeners("$podName-${Util.MAIN_QUEUE}", mainQueueListener)
            initializeListeners("$podName-${Util.DLQ}", dlqListener)
            LOG.info("Initialized SQS Queue Response Handler")
        }
    } catch (e: Exception) {
        LOG.warn("Failed to Initialize response handler due to ${e.message}")
    }
}
private fun initializeListeners(queueName: String, messageListener: MessageListener) {
    val connection = getSqsConnectionObject()
    val session = connection.createSession(false, SQSSession.UNORDERED_ACKNOWLEDGE)
    val consumer = session.createConsumer(session.createQueue(queueName))
    consumer.messageListener = messageListener
    connection.start()
    LOG.info("Initialized $messageListener for $queueName")
}

private fun getSqsConnectionObject(): SQSConnection {
    val connectionFactory = SQSConnectionFactory(
            ProviderConfiguration(),
            AmazonSQSClientBuilder.defaultClient()
    )
    return connectionFactory.createConnection()
}

DLQ监听器代码

@Named("dlqListener")
class DeadLetterQueueListener : MessageListener {

    @Value("\${pod_name}")
    private val podName: String = "dev"

    @Inject
    private lateinit var tradeStatusService: TradeStatusService

    companion object {
        private val LOG = LoggerFactory.getLogger(DeadLetterQueueListener::class.java)
    }

    override fun onMessage(message: Message?) {
        try {
            if (message is TextMessage) {
                LOG.info("Received message ${message.text} in Dead Letter Queue Listener " +
                        "with jmsMessageId: ${message.jmsMessageID}")
                if(message.getStringProperty(Util.QUEUE_NAME) == "$podName-${Util.MAIN_QUEUE}") {
                    message.acknowledge()
                    val jsonMessage = message.text
                    val objectMapper = ObjectMapper()
                    val externalTradeInputs = objectMapper.readValue(jsonMessage, ExternalTradeInputs::class.java)
                    val tradeStatusData = TradeStatusData().apply {
                        this.tradeBookingId = externalTradeInputs.tradeBookingId
                        this.statusDateTime = Date()
                    }
                }
            }
        } catch (e: Exception) {
            LOG.error("Error processing in Dead Letter Queue Listener with : ${e.message}")
        }
    }
}
排查参考提示
  • 验证队列名称完全匹配:确认代码中拼接的DLQ名称$podName-${Util.DLQ}与AWS控制台中实际DLQ的名称完全一致,包括大小写、前缀后缀、特殊字符等。同时检查pod_name环境变量的注入值是否正确,Util.DLQ常量的定义是否无误,避免因名称拼接错误导致监听器连接到了一个不存在(或自动创建的空队列)。
  • 检查IAM权限配置:确认应用程序使用的IAM角色拥有DLQ的sqs:ReceiveMessage、sqs:DeleteMessage、sqs:GetQueueAttributes等必要权限。主队列有权限不代表DLQ自动继承权限,需单独配置DLQ的权限策略。
  • 排查监听器初始化是否成功:当前init方法的catch块仅打印异常消息未输出堆栈,可能掩盖了DLQ监听器初始化失败的问题。建议修改日志逻辑,打印完整异常堆栈,或在initializeListeners方法内添加独立的异常捕获,确认DLQ的连接、会话、消费者是否真的成功创建并启动。
  • 确认JMS创建队列的行为:代码中使用session.createQueue(queueName)时,若指定队列不存在,AWS JMS客户端会自动创建新队列。检查AWS控制台是否存在名称近似的额外队列,确认监听器是否误连接到了自动创建的空队列,而非实际存储消息的DLQ。
  • 检查消息属性过滤逻辑:DLQ监听器中有message.getStringProperty(Util.QUEUE_NAME) == "$podName-${Util.MAIN_QUEUE}"的判断,若DLQ中的消息没有该属性或属性值不匹配,会导致消息收到后无后续处理,但日志会打印"Received message..."。若日志无输出则排除此情况,若有输出则需验证DLQ消息的属性值是否符合预期。
  • 配置长轮询减少空轮询:AWS SQS默认短轮询可能导致偶尔无法获取消息,可在创建消费者时设置ReceiveWaitTimeSeconds参数启用长轮询,或在DLQ的队列配置中开启长轮询,提升消息获取的可靠性。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 09:30:38