From 5d82b9d4a715749992da6929be6d70d7d9ea7163 Mon Sep 17 00:00:00 2001 From: lnk Date: Thu, 23 Jul 2026 15:47:07 +0800 Subject: [PATCH] =?UTF-8?q?=E6=B7=BB=E5=8A=A0=E7=94=9F=E4=BA=A7=E5=92=8C?= =?UTF-8?q?=E6=B6=88=E8=B4=B9=E6=8F=90=E7=A4=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- redisstream/RedisStreamMQ.cpp | 32 +++++++++++++++++++++++++++++--- 1 file changed, 29 insertions(+), 3 deletions(-) diff --git a/redisstream/RedisStreamMQ.cpp b/redisstream/RedisStreamMQ.cpp index 707ed9d..ebccb0f 100644 --- a/redisstream/RedisStreamMQ.cpp +++ b/redisstream/RedisStreamMQ.cpp @@ -686,6 +686,17 @@ private: } if (r->type == REDIS_REPLY_STRING) { + std::string message_id = r->str + ? std::string(r->str, r->len) + : std::string(); + std::cout << "[REDIS_STREAM][SEND_OK]" + << " stream=" << stream + << " id=" << message_id + << " key=" << key + << " tag=" << tag + << " enc=" << enc + << " body_len=" << body.size() + << std::endl; freeReplyObject(r); return REDIS_XADD_OK; } @@ -1023,13 +1034,19 @@ private: return NULL; } - void ack(const std::string& stream, const std::string& id) + bool ack(const std::string& stream, const std::string& id) { redisReply* r = conn_.command("XACK %b %b %b", stream.data(), stream.size(), group_.data(), group_.size(), id.data(), id.size()); - if (r) freeReplyObject(r); + if (!r) { + return false; + } + + bool acked = (r->type == REDIS_REPLY_INTEGER && r->integer > 0); + freeReplyObject(r); + return acked; } void handleOne(const Sub& sub, const std::string& id, redisReply* fieldsReply) @@ -1054,7 +1071,16 @@ private: rocketmq::ConsumeStatus ret = cb(msg); if (ret == rocketmq::CONSUME_SUCCESS) { - ack(sub.stream, id); + bool acked = ack(sub.stream, id); + std::cout << "[REDIS_STREAM][CONSUME_OK]" + << " stream=" << sub.stream + << " id=" << id + << " topic=" << sub.topic + << " tag=" << msg.getTags() + << " key=" << msg.getKeys() + << " body_len=" << msg.getBody().size() + << " ack=" << (acked ? "OK" : "FAIL") + << std::endl; } else { std::cout << "[REDIS_STREAM][RETRY_LATER] stream=" << sub.stream << " id=" << id << std::endl;