Получение копий каждой записи в SendMessageBatchRequest — AWS SQSJAVA

Программисты JAVA общаются здесь
Anonymous
Получение копий каждой записи в SendMessageBatchRequest — AWS SQS

Сообщение Anonymous »

Мы столкнулись с проблемой использования записей SQS. Мы отправляем пакетные записи в SQS из ECS, которые позже используются из Lambda. Я регистрирую каждую запись, отправляемую из ECS, и могу подтвердить, что это правильное количество записей, но по какой-то причине лямбда всегда также использует копию исходной записи. Я понятия не имею, откуда взялась эта копия; буквально через секунду в SQS всплывает оригинальный. У него такое же тело, но явно другой идентификатор сообщения.
Это код производителя:

Код: Выделить всё

public void sendMessageInBatch(List authorizedTransactions) {
List entries = new ArrayList();
for (int i = 0; i < authorizedTransactions.size(); i++) {
entries.add((SendMessageBatchRequestEntry.builder().id(String.valueOf(i+1)).messageBody(String.valueOf(authorizedTransactions.get(i))).build()));
if ((i+1) % 10 == 0) {
this.sendMessagesBatch(entries);
entries = new ArrayList();
}
}
if (!entries.isEmpty()) {
this.sendMessagesBatch(entries);
}
}

@Async
private void sendMessagesBatch(List entries) {
List successfulEntries;
List failedEntries;
entries.forEach((entry) -> {
LOG.info("ID: " + entry.id() + " Transaction: " + entry.messageBody());
});
try {
SendMessageBatchRequest sendMessageBatchRequest = SendMessageBatchRequest.builder().queueUrl(queueUrl).entries(entries).build();
SendMessageBatchResponse sendMessageBatchResponse = sqsClient.sendMessageBatch(sendMessageBatchRequest);
if (sendMessageBatchResponse.hasSuccessful()) {
successfulEntries = sendMessageBatchResponse.successful();
LOG.info("Processed below entries: ");
successfulEntries.forEach((entry) -> {
LOG.info(entry.messageId());
});
}
if (sendMessageBatchResponse.hasFailed()) {
failedEntries = sendMessageBatchResponse.failed();
LOG.info("Below entries failed: ");
failedEntries.forEach((entry) -> {
LOG.info(entry.id() + " " + entry.message());
});
}
} catch (Exception exception) {
throw new RuntimeException(exception);
}
}

Это инфраструктурный код:

Код: Выделить всё

this.pollingQueue = new Queue(this, 'auth-record-polling-queue', {
queueName: 'auth-record-polling-queue',
retentionPeriod: Duration.days(1),
visibilityTimeout: Duration.minutes(10),
deliveryDelay: Duration.seconds(15),
deadLetterQueue: {
queue: this.deadLetterQueue,
maxReceiveCount: 2
}
});

Поскольку потребительский код представляет собой лямбду и слишком велик, чтобы его можно было здесь приводить, это репозиторий кода GitHub и лямбда-код: Lambda Handler
Я использую идемпотентная оболочка вокруг обработчика... Я не уверен, что это вызывает проблему.
Я пробовал следующее:
Увеличение тайм-аут видимости SQS
Производители журналов отправляют
Добавление задержки перед потреблением.

Подробнее здесь: https://stackoverflow.com/questions/790 ... st-aws-sqs

Вернуться в «JAVA»