From 1406f246a3199e6d0832bb4e01ac13ef9d3c06c5 Mon Sep 17 00:00:00 2001 From: TB Date: Thu, 16 Jul 2026 17:33:27 +0800 Subject: [PATCH 1/7] compatiable ror bitmap rdb type --- .../core/redis/operation/RedisOpType.java | 2 ++ .../redis/core/redis/rdb/RdbConstant.java | 2 +- .../redis/core/redis/rdb/RdbParseContext.java | 2 +- .../rdb/parser/DefaultRdbParseContext.java | 3 ++ .../redis/rdb/parser/RdbBitmapParser.java | 7 +++++ .../rdb/parser/DefaultRdbParserTest.java | 29 +++++++++++++++++-- .../core/redis/rdb/parser/RdbDataBytes.java | 3 ++ 7 files changed, 44 insertions(+), 4 deletions(-) diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/RedisOpType.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/RedisOpType.java index b76baadf24..20e42846c7 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/RedisOpType.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/RedisOpType.java @@ -83,6 +83,8 @@ public enum RedisOpType { // Bit single SETBIT(false, 4), + BITCOUNT(false, 4), + // String multi DEL(true, -2), diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbConstant.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbConstant.java index 8050151897..40fda6681a 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbConstant.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbConstant.java @@ -40,7 +40,7 @@ private RdbConstant() { public static final short REDIS_RDB_TYPE_HASH_ZIPLIST = 13; public static final short REDIS_RDB_TYPE_LIST_QUICKLIST = 14; public static final short REDIS_RDB_TYPE_STREAM_LISTPACKS = 15; -// public static final short REDIS_RDB_TYPE_BITMAP = 16; + public static final short REDIS_RDB_TYPE_BITMAP = 192; public static final short REDIS_RDB_TYPE_CRDT = 200; // Redis 8.x 新增操作码 (rdb.h) diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbParseContext.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbParseContext.java index c335bd1705..e9a74f42cd 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbParseContext.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbParseContext.java @@ -90,7 +90,7 @@ enum RdbType { HASH_ZIPLIST(RdbConstant.REDIS_RDB_TYPE_HASH_ZIPLIST, false, RdbHashZipListParser::new), LIST_QUICKLIST(RdbConstant.REDIS_RDB_TYPE_LIST_QUICKLIST, false, RdbQuickListParser::new), STREAM_LISTPACKS(RdbConstant.REDIS_RDB_TYPE_STREAM_LISTPACKS, false, RdbStreamListpacksParser::new), -// BITMAP(RdbConstant.REDIS_RDB_TYPE_BITMAP, false, RdbBitmapParser::new), + BITMAP(RdbConstant.REDIS_RDB_TYPE_BITMAP, false, RdbBitmapParser::new), CRDT(RdbConstant.REDIS_RDB_TYPE_CRDT, false, DefaultRdbCrdtParser::new), // MODULE_AUX(RdbConstant.REDIS_RDB_OP_CODE_MODULE_AUX), IDLE(RdbConstant.REDIS_RDB_OP_CODE_IDLE, true, RdbIdleParser::new), diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParseContext.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParseContext.java index e6a469cb4b..61d98ffee0 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParseContext.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParseContext.java @@ -51,6 +51,9 @@ public synchronized void bindRdbParser(RdbParser parser) { @Override public RdbParser getOrCreateParser(RdbType rdbType) { + if(this.getRdbVersion() <= 9 && rdbType.getCode() == 16){ + rdbType = RdbType.BITMAP; + } RdbParser parser = parsers.get(rdbType); if (null != parser) return parser; diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/RdbBitmapParser.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/RdbBitmapParser.java index 5fb4618142..afde743eac 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/RdbBitmapParser.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/RdbBitmapParser.java @@ -114,6 +114,13 @@ private void propagateCmdIfNeed(byte[] value) { RedisOpType.SET, new byte[][]{RedisOpType.SET.name().getBytes(), context.getKey().get(), value}, context.getKey(), value)); + + if(this.context.getRdbVersion() > 9) { + notifyRedisOp(new RedisOpSingleKey( + RedisOpType.BITCOUNT, + new byte[][]{RedisOpType.BITCOUNT.name().getBytes(), context.getKey().get()}, + context.getKey(), value)); + } } @Override diff --git a/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParserTest.java b/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParserTest.java index 7efa4534b4..8daf473502 100644 --- a/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParserTest.java +++ b/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParserTest.java @@ -326,7 +326,6 @@ public void testParseCrdtSortedSet() { Assert.assertEquals("ZADD testSS329 0.2700072245105105 4key89", redisOps.get(19).toString()); } - /* @Test public void testParseBitmap() { ByteBuf byteBuf = Unpooled.wrappedBuffer(rorBitmap); @@ -347,7 +346,33 @@ public void testParseBitmap() { Assert.assertNotEquals(0, bitmapKey[2][8] & (1 << 1)); // getbit bitmap_key 70 -> 1 Assert.assertEquals("SET common_key value", redisOps.get(4).toString()); } - */ + + @Test + public void testParseRor8Bitmap() { + ByteBuf byteBuf = Unpooled.wrappedBuffer(ror8Bitmap); + while (!parser.isFinish()) { + parser.read(byteBuf); + } + + Assert.assertEquals("SELECT 0", redisOps.get(0).toString()); + + byte[][] bitmapKey2 = redisOps.get(1).buildRawOpArgs(); + Assert.assertEquals("bitmap_key2", new String(bitmapKey2[1])); + Assert.assertEquals(700 / 8 + (700 % 8 == 0 ? 0 : 1),bitmapKey2[2].length); + Assert.assertEquals("BITCOUNT bitmap_key2", redisOps.get(2).toString()); + + Assert.assertEquals("SET common_key value", redisOps.get(3).toString()); + + Assert.assertEquals("SET mykey aa", redisOps.get(5).toString()); + + byte[][] bitmapKey = redisOps.get(7).buildRawOpArgs(); + Assert.assertEquals("bitmap_key", new String(bitmapKey[1])); + Assert.assertEquals(700 / 8 + (700 % 8 == 0 ? 0 : 1),bitmapKey[2].length); + Assert.assertNotEquals(0, bitmapKey[2][87] & (1 << 3)); // getbit bitmap_key 700 -> 1 + Assert.assertNotEquals(0, bitmapKey[2][8] & (1 << 1)); // getbit bitmap_key 70 -> 1 + + Assert.assertEquals(9,redisOps.size()); + } @Test public void testParseListpackStream3() { diff --git a/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/RdbDataBytes.java b/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/RdbDataBytes.java index f403c4f122..6b491f36a9 100644 --- a/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/RdbDataBytes.java +++ b/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/RdbDataBytes.java @@ -430,6 +430,9 @@ public class RdbDataBytes { 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, (byte) 0x08, (byte) 0xf9, 0x05, 0x10, (byte) 0x0a, 0x63, (byte) 0x6f, (byte) 0x6d, (byte) 0x6d, (byte) 0x6f, (byte) 0x6e, (byte) 0x5f, (byte) 0x6b, 0x65, 0x79, 0x05, 0x50, 0x00, 0x05, 0x76, 0x61, (byte) 0x6c, 0x75, 0x65, (byte) 0xff, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00}; + public static final byte[] ror8Bitmap = new byte[]{0x52, 0x45, 0x44, 0x49, 0x53, 0x30, 0x30, 0x31, 0x32, (byte) 0xFA, 0x09, 0x72, 0x65, 0x64, 0x69, 0x73, 0x2D, 0x76, 0x65, 0x72, 0x05, 0x38, 0x2E, 0x32, 0x2E, 0x30, (byte) 0xFA, 0x0A, 0x72, 0x65, 0x64, 0x69, 0x73, 0x2D, 0x62, 0x69, 0x74, 0x73, (byte) 0xC0, 0x40, (byte) 0xFA, 0x05, 0x63, 0x74, 0x69, 0x6D, 0x65, (byte) 0xC2, 0x44, (byte) 0x81, 0x58, 0x6A, (byte) 0xFA, 0x08, 0x75, 0x73, 0x65, 0x64, 0x2D, 0x6D, 0x65, 0x6D, (byte) 0xC2, 0x78, (byte) 0x9E, (byte) 0x97, 0x02, (byte) 0xFA, 0x0E, 0x72, 0x65, 0x70, 0x6C, 0x2D, 0x73, 0x74, 0x72, 0x65, 0x61, 0x6D, 0x2D, 0x64, 0x62, (byte) 0xC0, 0x00, (byte) 0xFA, 0x07, 0x72, 0x65, 0x70, 0x6C, 0x2D, 0x69, 0x64, 0x28, 0x30, 0x32, 0x32, 0x32, 0x36, 0x31, 0x66, 0x37, 0x31, 0x65, 0x36, 0x66, 0x35, 0x65, 0x33, 0x65, 0x61, 0x62, 0x32, 0x64, 0x62, 0x33, 0x33, 0x63, 0x37, 0x32, 0x39, 0x34, 0x39, 0x30, 0x30, 0x39, 0x34, 0x63, 0x30, 0x37, 0x36, 0x33, 0x37, 0x34, (byte) 0xFA, 0x0B, 0x72, 0x65, 0x70, 0x6C, 0x2D, 0x6F, 0x66, 0x66, 0x73, 0x65, 0x74, (byte) 0xC1, (byte) 0xA9, 0x00, (byte) 0xFA, 0x0E, 0x67, 0x74, 0x69, 0x64, 0x2D, 0x72, 0x65, 0x70, 0x6C, 0x2D, 0x6D, 0x6F, 0x64, 0x65, 0x05, 0x70, 0x73, 0x79, 0x6E, 0x63, (byte) 0xFA, 0x0D, 0x67, 0x74, 0x69, 0x64, 0x2D, 0x65, 0x78, 0x65, 0x63, 0x75, 0x74, 0x65, 0x64, 0x00, (byte) 0xFA, 0x09, 0x67, 0x74, 0x69, 0x64, 0x2D, 0x6C, 0x6F, 0x73, 0x74, 0x00, (byte) 0xFA, 0x08, 0x61, 0x6F, 0x66, 0x2D, 0x62, 0x61, 0x73, 0x65, (byte) 0xC0, 0x00, (byte) 0xFE, 0x00, (byte) 0xFB, 0x04, 0x00, (byte) 0xF9, 0x00, (byte) 0xC0, 0x0B, 0x62, 0x69, 0x74, 0x6D, 0x61, 0x70, 0x5F, 0x6B, 0x65, 0x79, 0x32, 0x40, 0x58, 0x50, 0x00, 0x40, 0x58, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x08, (byte) 0xF9, 0x00, (byte) 0xC0, 0x0A, 0x63, 0x6F, 0x6D, 0x6D, 0x6F, 0x6E, 0x5F, 0x6B, 0x65, 0x79, 0x05, 0x50, 0x00, 0x05, 0x76, 0x61, 0x6C, 0x75, 0x65, (byte) 0xF9, 0x00, (byte) 0xC0, 0x05, 0x6D, 0x79, 0x6B, 0x65, 0x79, 0x02, 0x50, 0x00, 0x02, 0x61, 0x61, (byte) 0xF9, 0x00, (byte) 0xC0, 0x0A, 0x62, 0x69, 0x74, 0x6D, 0x61, 0x70, 0x5F, 0x6B, 0x65, 0x79, 0x40, 0x58, 0x50, 0x00, 0x40, 0x58, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x08, (byte) 0xFF, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00}; + + public static final byte[] listpackStream3 = new byte[]{0x52, 0x45, 0x44, 0x49, 0x53, 0x30, 0x30, 0x31, 0x32, (byte) 0xFA, 0x09, 0x72, 0x65, 0x64, 0x69, 0x73, 0x2D, 0x76, 0x65, 0x72, 0x05, 0x38, 0x2E, 0x32, 0x2E, 0x30, (byte) 0xFA, 0x0A, 0x72, 0x65, 0x64, 0x69, 0x73, 0x2D, 0x62, 0x69, 0x74, 0x73, (byte) 0xC0, 0x40, (byte) 0xFA, 0x05, 0x63, 0x74, 0x69, 0x6D, 0x65, (byte) 0xC2, (byte) 0xEC, (byte) 0xFB, (byte) 0xC0, 0x69, (byte) 0xFA, 0x08, 0x75, 0x73, 0x65, 0x64, 0x2D, 0x6D, 0x65, 0x6D, (byte) 0xC2, (byte) 0xD0, 0x18, 0x13, 0x00, (byte) 0xFA, 0x08, 0x61, 0x6F, 0x66, 0x2D, 0x62, 0x61, 0x73, 0x65, (byte) 0xC0, 0x00, (byte) 0xFE, 0x00, (byte) 0xFB, 0x01, 0x00, 0x15, 0x08, 0x6D, 0x79, 0x73, 0x74, 0x72, 0x65, 0x61, 0x6D, 0x01, 0x10, 0x00, 0x00, 0x01, (byte) 0x9D, 0x19, (byte) 0xCE, (byte) 0xED, (byte) 0xD2, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, (byte) 0xC3, 0x40, 0x61, 0x40, 0x71, 0x06, 0x71, 0x00, 0x00, 0x00, 0x17, 0x00, 0x01, 0x20, 0x00, 0x0B, 0x02, 0x01, (byte) 0x84, 0x6E, 0x61, 0x6D, 0x65, 0x05, (byte) 0x87, 0x73, 0x75, 0x72, 0x40, 0x08, 0x05, 0x08, 0x00, 0x01, 0x03, 0x01, 0x00, 0x20, 0x01, 0x0F, (byte) 0x84, 0x53, 0x61, 0x72, 0x61, 0x05, (byte) 0x87, 0x4F, 0x43, 0x6F, 0x6E, 0x6E, 0x6F, 0x72, 0x08, 0x05, 0x20, 0x12, 0x03, (byte) 0xF1, (byte) 0xD9, 0x3F, 0x03, 0x40, 0x1E, 0x0D, (byte) 0x86, 0x66, 0x69, 0x65, 0x6C, 0x64, 0x31, 0x07, (byte) 0x86, 0x76, 0x61, 0x6C, 0x75, 0x65, 0x20, 0x07, 0x60, 0x0F, 0x00, 0x32, (byte) 0xA0, 0x0F, 0x20, 0x07, 0x60, 0x0F, 0x00, 0x33, (byte) 0xA0, 0x0F, 0x04, 0x33, 0x07, 0x0A, 0x01, (byte) 0xFF, 0x01, (byte) 0x81, 0x00, 0x00, 0x01, (byte) 0x9D, 0x19, (byte) 0xCF, 0x2D, (byte) 0xAB, 0x00, (byte) 0x81, 0x00, 0x00, 0x01, (byte) 0x9D, 0x19, (byte) 0xCF, 0x2D, (byte) 0xAB, 0x00, (byte) 0x81, 0x00, 0x00, 0x01, (byte) 0x9D, 0x19, (byte) 0xCE, (byte) 0xED, (byte) 0xD2, 0x00, 0x02, 0x01, 0x07, 0x6D, 0x79, 0x67, 0x72, 0x6F, 0x75, 0x70, (byte) 0x81, 0x00, 0x00, 0x01, (byte) 0x9D, 0x19, (byte) 0xCF, 0x2D, (byte) 0xAB, 0x00, 0x02, 0x01, 0x00, 0x00, 0x01, (byte) 0x9D, 0x19, (byte) 0xCF, 0x2D, (byte) 0xAB, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x21, 0x7C, (byte) 0xD3, 0x19, (byte) 0x9D, 0x01, 0x00, 0x00, 0x01, 0x01, 0x0A, 0x6D, 0x79, 0x63, 0x6F, 0x6E, 0x73, 0x75, 0x6D, 0x65, 0x72, 0x21, 0x7C, (byte) 0xD3, 0x19, (byte) 0x9D, 0x01, 0x00, 0x00, 0x21, 0x7C, (byte) 0xD3, 0x19, (byte) 0x9D, 0x01, 0x00, 0x00, 0x01, 0x00, 0x00, 0x01, (byte) 0x9D, 0x19, (byte) 0xCF, 0x2D, (byte) 0xAB, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, (byte) 0xFF, (byte) 0xAE, (byte) 0xBE, 0x77, 0x6F, 0x5D, (byte) 0xBC, (byte) 0x9C, (byte) 0x8C}; public static final byte[] listpackHashEx = new byte[]{0x52, 0x45, 0x44, 0x49, 0x53, 0x30, 0x30, 0x31, 0x32, (byte) 0xFA, 0x09, 0x72, 0x65, 0x64, 0x69, 0x73, 0x2D, 0x76, 0x65, 0x72, 0x05, 0x38, 0x2E, 0x32, 0x2E, 0x30, (byte) 0xFA, 0x0A, 0x72, 0x65, 0x64, 0x69, 0x73, 0x2D, 0x62, 0x69, 0x74, 0x73, (byte) 0xC0, 0x40, (byte) 0xFA, 0x05, 0x63, 0x74, 0x69, 0x6D, 0x65, (byte) 0xC2, (byte) 0x94, 0x25, (byte) 0xC1, 0x69, (byte) 0xFA, 0x08, 0x75, 0x73, 0x65, 0x64, 0x2D, 0x6D, 0x65, 0x6D, (byte) 0xC2, 0x30, (byte) 0xE5, 0x14, 0x00, (byte) 0xFA, 0x08, 0x61, 0x6F, 0x66, 0x2D, 0x62, 0x61, 0x73, 0x65, (byte) 0xC0, 0x00, (byte) 0xFE, 0x00, (byte) 0xFB, 0x01, 0x00, 0x19, 0x07, 0x6D, 0x79, 0x65, 0x78, 0x6B, 0x65, 0x79, (byte) 0xA2, (byte) 0xA4, (byte) 0xEE, 0x2F, (byte) 0x9D, 0x01, 0x00, 0x00, (byte) 0xC3, 0x2F, 0x3D, 0x14, 0x3D, 0x00, 0x00, 0x00, 0x09, 0x00, (byte) 0x82, 0x66, 0x33, 0x03, (byte) 0x82, 0x76, 0x33, 0x03, (byte) 0xF4, (byte) 0xA2, (byte) 0xA4, (byte) 0xEE, 0x2F, (byte) 0x9D, 0x01, 0x20, 0x12, 0x02, (byte) 0x82, 0x66, 0x32, 0x20, 0x11, 0x00, 0x32, (byte) 0xE0, 0x04, 0x11, 0x00, 0x31, 0x20, 0x11, 0x00, 0x31, (byte) 0xE0, 0x01, 0x11, 0x01, 0x09, (byte) 0xFF, (byte) 0xFF, 0x5D, (byte) 0x89, (byte) 0xB4, 0x4E, (byte) 0xBA, (byte) 0xC3, (byte) 0xDF, (byte) 0xBA}; From 04a3858ded41657dc0d2036a9492392b63037a01 Mon Sep 17 00:00:00 2001 From: TB Date: Mon, 20 Jul 2026 13:05:12 +0800 Subject: [PATCH 2/7] rdb type add version --- .../redis/core/redis/rdb/RdbConstant.java | 1 + .../redis/core/redis/rdb/RdbParseContext.java | 60 +++++++++++++------ .../rdb/parser/DefaultRdbParseContext.java | 3 - .../redis/rdb/parser/DefaultRdbParser.java | 2 +- 4 files changed, 45 insertions(+), 21 deletions(-) diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbConstant.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbConstant.java index 40fda6681a..b3acf870bb 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbConstant.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbConstant.java @@ -41,6 +41,7 @@ private RdbConstant() { public static final short REDIS_RDB_TYPE_LIST_QUICKLIST = 14; public static final short REDIS_RDB_TYPE_STREAM_LISTPACKS = 15; public static final short REDIS_RDB_TYPE_BITMAP = 192; + public static final short REDIS_RDB_TYPE_BITMAP_9 = 16; public static final short REDIS_RDB_TYPE_CRDT = 200; // Redis 8.x 新增操作码 (rdb.h) diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbParseContext.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbParseContext.java index e9a74f42cd..4a5454dcb2 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbParseContext.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbParseContext.java @@ -90,7 +90,8 @@ enum RdbType { HASH_ZIPLIST(RdbConstant.REDIS_RDB_TYPE_HASH_ZIPLIST, false, RdbHashZipListParser::new), LIST_QUICKLIST(RdbConstant.REDIS_RDB_TYPE_LIST_QUICKLIST, false, RdbQuickListParser::new), STREAM_LISTPACKS(RdbConstant.REDIS_RDB_TYPE_STREAM_LISTPACKS, false, RdbStreamListpacksParser::new), - BITMAP(RdbConstant.REDIS_RDB_TYPE_BITMAP, false, RdbBitmapParser::new), + BITMAP_12(12,RdbConstant.REDIS_RDB_TYPE_BITMAP, false, RdbBitmapParser::new), + BITMAP(RdbConstant.REDIS_RDB_TYPE_BITMAP_9, false, RdbBitmapParser::new), CRDT(RdbConstant.REDIS_RDB_TYPE_CRDT, false, DefaultRdbCrdtParser::new), // MODULE_AUX(RdbConstant.REDIS_RDB_OP_CODE_MODULE_AUX), IDLE(RdbConstant.REDIS_RDB_OP_CODE_IDLE, true, RdbIdleParser::new), @@ -116,22 +117,23 @@ enum RdbType { // Redis 8.x 新增 - KEY_META(RdbConstant.RDB_OPCODE_KEY_META, true, RdbKeyMetaParser::new), - FUNCTION2(RdbConstant.RDB_OPCODE_FUNCTION2, true, RdbFunction2Parser::new), - MODULE_AUX(RdbConstant.REDIS_RDB_OP_CODE_MODULE_AUX, true, RdbModuleAuxParser::new), - MODULE_2(RdbConstant.REDIS_RDB_TYPE_MODULE_2, true, RdbModuleParser::new), - - HASH_LISTPACK(RdbConstant.REDIS_RDB_TYPE_HASH_LISTPACK, false, RdbHashListpackParser::new), - ZSET_LISTPACK(RdbConstant.REDIS_RDB_TYPE_ZSET_LISTPACK, false, RdbZSetListpackParser::new), - - LIST_QUICKLIST_2(RdbConstant.REDIS_RDB_TYPE_LIST_QUICKLIST_2, false, RdbQuickList2Parser::new), - HASH_METADATA(RdbConstant.REDIS_RDB_TYPE_HASH_METADATA,false,RdbHashMetadataParser::new), - HASH_LISTPACK_EX(RdbConstant.REDIS_RDB_TYPE_HASH_LISTPACK_EX, false, RdbHashListpackExParser::new), - STREAM_LISTPACKS_2(RdbConstant.REDIS_RDB_TYPE_STREAM_LISTPACKS_2, false, RdbStreamListpacks2Parser::new), - SET_LISTPACK(RdbConstant.REDIS_RDB_TYPE_SET_LISTPACK, false, RdbSetListpackParser::new), - STREAM_LISTPACKS_3(RdbConstant.REDIS_RDB_TYPE_STREAM_LISTPACKS_3, false, RdbStreamListpacks3Parser::new), - STREAM_LISTPACKS_4(RdbConstant.REDIS_RDB_TYPE_STREAM_LISTPACKS_4, false, RdbStreamListpacks4Parser::new); - + KEY_META(12,RdbConstant.RDB_OPCODE_KEY_META, true, RdbKeyMetaParser::new), + FUNCTION2(12,RdbConstant.RDB_OPCODE_FUNCTION2, true, RdbFunction2Parser::new), + MODULE_AUX(12,RdbConstant.REDIS_RDB_OP_CODE_MODULE_AUX, true, RdbModuleAuxParser::new), + MODULE_2(12,RdbConstant.REDIS_RDB_TYPE_MODULE_2, true, RdbModuleParser::new), + + HASH_LISTPACK(12,RdbConstant.REDIS_RDB_TYPE_HASH_LISTPACK, false, RdbHashListpackParser::new), + ZSET_LISTPACK(12,RdbConstant.REDIS_RDB_TYPE_ZSET_LISTPACK, false, RdbZSetListpackParser::new), + + LIST_QUICKLIST_2(12,RdbConstant.REDIS_RDB_TYPE_LIST_QUICKLIST_2, false, RdbQuickList2Parser::new), + HASH_METADATA(12,RdbConstant.REDIS_RDB_TYPE_HASH_METADATA,false,RdbHashMetadataParser::new), + HASH_LISTPACK_EX(12,RdbConstant.REDIS_RDB_TYPE_HASH_LISTPACK_EX, false, RdbHashListpackExParser::new), + STREAM_LISTPACKS_2(12,RdbConstant.REDIS_RDB_TYPE_STREAM_LISTPACKS_2, false, RdbStreamListpacks2Parser::new), + SET_LISTPACK(12,RdbConstant.REDIS_RDB_TYPE_SET_LISTPACK, false, RdbSetListpackParser::new), + STREAM_LISTPACKS_3(12,RdbConstant.REDIS_RDB_TYPE_STREAM_LISTPACKS_3, false, RdbStreamListpacks3Parser::new), + STREAM_LISTPACKS_4(12,RdbConstant.REDIS_RDB_TYPE_STREAM_LISTPACKS_4, false, RdbStreamListpacks4Parser::new); + + private int version; private short code; private boolean rdbOp; @@ -139,8 +141,17 @@ enum RdbType { private Function parserConstructor; private static final Map types = new HashMap<>(); + private static final Map> versionTypes = new HashMap<>(); RdbType(short code, boolean rdbOp, Function parserConstructor) { + this.version = 9; + this.code = code; + this.rdbOp = rdbOp; + this.parserConstructor = parserConstructor; + } + + RdbType(int version,short code, boolean rdbOp, Function parserConstructor) { + this.version = version; this.code = code; this.rdbOp = rdbOp; this.parserConstructor = parserConstructor; @@ -163,14 +174,29 @@ public RdbParser makeParser(RdbParseContext parserManager) { static { for (RdbType rdbType : values()) { + if(rdbType.version > 9) continue; types.put(rdbType.code, rdbType); } + for (RdbType rdbType : values()) { + if(rdbType.version <= 9) continue; + Map versionMap = versionTypes.computeIfAbsent(rdbType.version,(version)-> new HashMap<>()); + versionMap.put(rdbType.code,rdbType); + } } public static RdbType findByCode(short code) { return types.get(code); } + public static RdbType findByCode(int version,short code) { + Map versionMap = versionTypes.getOrDefault(version,types); + RdbType rdbType = versionMap.get(code); + if(rdbType == null){ + return findByCode(code); + } + return rdbType; + } + } String cSet = "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789-_"; diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParseContext.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParseContext.java index 61d98ffee0..e6a469cb4b 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParseContext.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParseContext.java @@ -51,9 +51,6 @@ public synchronized void bindRdbParser(RdbParser parser) { @Override public RdbParser getOrCreateParser(RdbType rdbType) { - if(this.getRdbVersion() <= 9 && rdbType.getCode() == 16){ - rdbType = RdbType.BITMAP; - } RdbParser parser = parsers.get(rdbType); if (null != parser) return parser; diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParser.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParser.java index 98c64f5235..eafda569d1 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParser.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParser.java @@ -89,7 +89,7 @@ public Void read(ByteBuf byteBuf) { case READ_TYPE: short type = byteBuf.readUnsignedByte(); - RdbParseContext.RdbType newType = RdbParseContext.RdbType.findByCode(type); + RdbParseContext.RdbType newType = RdbParseContext.RdbType.findByCode(rdbVersion,type); if (currentType == RdbParseContext.RdbType.AUX && newType != RdbParseContext.RdbType.AUX) { auxFinished.set(true); notifyAuxEnd(rdbParseContext.getAllAux()); From beb548b9df8060260a8c122d0d3ae0d703b4cc65 Mon Sep 17 00:00:00 2001 From: TB Date: Mon, 20 Jul 2026 15:38:39 +0800 Subject: [PATCH 3/7] add conflict type into version map --- .../redis/core/redis/rdb/RdbParseContext.java | 49 +++++++++---------- .../redis/rdb/parser/RdbBitmapParser.java | 10 ++-- .../rdb/parser/DefaultRdbParserTest.java | 6 +-- 3 files changed, 31 insertions(+), 34 deletions(-) diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbParseContext.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbParseContext.java index 4a5454dcb2..70bea04705 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbParseContext.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbParseContext.java @@ -90,8 +90,8 @@ enum RdbType { HASH_ZIPLIST(RdbConstant.REDIS_RDB_TYPE_HASH_ZIPLIST, false, RdbHashZipListParser::new), LIST_QUICKLIST(RdbConstant.REDIS_RDB_TYPE_LIST_QUICKLIST, false, RdbQuickListParser::new), STREAM_LISTPACKS(RdbConstant.REDIS_RDB_TYPE_STREAM_LISTPACKS, false, RdbStreamListpacksParser::new), - BITMAP_12(12,RdbConstant.REDIS_RDB_TYPE_BITMAP, false, RdbBitmapParser::new), - BITMAP(RdbConstant.REDIS_RDB_TYPE_BITMAP_9, false, RdbBitmapParser::new), + BITMAP(RdbConstant.REDIS_RDB_TYPE_BITMAP, false, RdbBitmapParser::new), + BITMAP_9(9,RdbConstant.REDIS_RDB_TYPE_BITMAP_9, false, RdbBitmapParser::new), CRDT(RdbConstant.REDIS_RDB_TYPE_CRDT, false, DefaultRdbCrdtParser::new), // MODULE_AUX(RdbConstant.REDIS_RDB_OP_CODE_MODULE_AUX), IDLE(RdbConstant.REDIS_RDB_OP_CODE_IDLE, true, RdbIdleParser::new), @@ -117,21 +117,21 @@ enum RdbType { // Redis 8.x 新增 - KEY_META(12,RdbConstant.RDB_OPCODE_KEY_META, true, RdbKeyMetaParser::new), - FUNCTION2(12,RdbConstant.RDB_OPCODE_FUNCTION2, true, RdbFunction2Parser::new), - MODULE_AUX(12,RdbConstant.REDIS_RDB_OP_CODE_MODULE_AUX, true, RdbModuleAuxParser::new), - MODULE_2(12,RdbConstant.REDIS_RDB_TYPE_MODULE_2, true, RdbModuleParser::new), - - HASH_LISTPACK(12,RdbConstant.REDIS_RDB_TYPE_HASH_LISTPACK, false, RdbHashListpackParser::new), - ZSET_LISTPACK(12,RdbConstant.REDIS_RDB_TYPE_ZSET_LISTPACK, false, RdbZSetListpackParser::new), - - LIST_QUICKLIST_2(12,RdbConstant.REDIS_RDB_TYPE_LIST_QUICKLIST_2, false, RdbQuickList2Parser::new), - HASH_METADATA(12,RdbConstant.REDIS_RDB_TYPE_HASH_METADATA,false,RdbHashMetadataParser::new), - HASH_LISTPACK_EX(12,RdbConstant.REDIS_RDB_TYPE_HASH_LISTPACK_EX, false, RdbHashListpackExParser::new), - STREAM_LISTPACKS_2(12,RdbConstant.REDIS_RDB_TYPE_STREAM_LISTPACKS_2, false, RdbStreamListpacks2Parser::new), - SET_LISTPACK(12,RdbConstant.REDIS_RDB_TYPE_SET_LISTPACK, false, RdbSetListpackParser::new), - STREAM_LISTPACKS_3(12,RdbConstant.REDIS_RDB_TYPE_STREAM_LISTPACKS_3, false, RdbStreamListpacks3Parser::new), - STREAM_LISTPACKS_4(12,RdbConstant.REDIS_RDB_TYPE_STREAM_LISTPACKS_4, false, RdbStreamListpacks4Parser::new); + KEY_META(RdbConstant.RDB_OPCODE_KEY_META, true, RdbKeyMetaParser::new), + FUNCTION2(RdbConstant.RDB_OPCODE_FUNCTION2, true, RdbFunction2Parser::new), + MODULE_AUX(RdbConstant.REDIS_RDB_OP_CODE_MODULE_AUX, true, RdbModuleAuxParser::new), + MODULE_2(RdbConstant.REDIS_RDB_TYPE_MODULE_2, true, RdbModuleParser::new), + + HASH_LISTPACK(RdbConstant.REDIS_RDB_TYPE_HASH_LISTPACK, false, RdbHashListpackParser::new), + ZSET_LISTPACK(RdbConstant.REDIS_RDB_TYPE_ZSET_LISTPACK, false, RdbZSetListpackParser::new), + + LIST_QUICKLIST_2(RdbConstant.REDIS_RDB_TYPE_LIST_QUICKLIST_2, false, RdbQuickList2Parser::new), + HASH_METADATA(RdbConstant.REDIS_RDB_TYPE_HASH_METADATA,false,RdbHashMetadataParser::new), + HASH_LISTPACK_EX(RdbConstant.REDIS_RDB_TYPE_HASH_LISTPACK_EX, false, RdbHashListpackExParser::new), + STREAM_LISTPACKS_2(RdbConstant.REDIS_RDB_TYPE_STREAM_LISTPACKS_2, false, RdbStreamListpacks2Parser::new), + SET_LISTPACK(RdbConstant.REDIS_RDB_TYPE_SET_LISTPACK, false, RdbSetListpackParser::new), + STREAM_LISTPACKS_3(RdbConstant.REDIS_RDB_TYPE_STREAM_LISTPACKS_3, false, RdbStreamListpacks3Parser::new), + STREAM_LISTPACKS_4(RdbConstant.REDIS_RDB_TYPE_STREAM_LISTPACKS_4, false, RdbStreamListpacks4Parser::new); private int version; private short code; @@ -144,7 +144,7 @@ enum RdbType { private static final Map> versionTypes = new HashMap<>(); RdbType(short code, boolean rdbOp, Function parserConstructor) { - this.version = 9; + this.version = 0; this.code = code; this.rdbOp = rdbOp; this.parserConstructor = parserConstructor; @@ -174,13 +174,12 @@ public RdbParser makeParser(RdbParseContext parserManager) { static { for (RdbType rdbType : values()) { - if(rdbType.version > 9) continue; - types.put(rdbType.code, rdbType); - } - for (RdbType rdbType : values()) { - if(rdbType.version <= 9) continue; - Map versionMap = versionTypes.computeIfAbsent(rdbType.version,(version)-> new HashMap<>()); - versionMap.put(rdbType.code,rdbType); + if(rdbType.version != 0) { + Map versionMap = versionTypes.computeIfAbsent(rdbType.version, (version) -> new HashMap<>()); + versionMap.put(rdbType.code, rdbType); + }else { + types.put(rdbType.code, rdbType); + } } } diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/RdbBitmapParser.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/RdbBitmapParser.java index afde743eac..7d957bffdd 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/RdbBitmapParser.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/RdbBitmapParser.java @@ -115,12 +115,10 @@ private void propagateCmdIfNeed(byte[] value) { new byte[][]{RedisOpType.SET.name().getBytes(), context.getKey().get(), value}, context.getKey(), value)); - if(this.context.getRdbVersion() > 9) { - notifyRedisOp(new RedisOpSingleKey( - RedisOpType.BITCOUNT, - new byte[][]{RedisOpType.BITCOUNT.name().getBytes(), context.getKey().get()}, - context.getKey(), value)); - } + notifyRedisOp(new RedisOpSingleKey( + RedisOpType.BITCOUNT, + new byte[][]{RedisOpType.BITCOUNT.name().getBytes(), context.getKey().get()}, + context.getKey(), value)); } @Override diff --git a/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParserTest.java b/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParserTest.java index 8daf473502..336355f9e5 100644 --- a/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParserTest.java +++ b/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/DefaultRdbParserTest.java @@ -336,15 +336,15 @@ public void testParseBitmap() { Assert.assertEquals("SELECT 0", redisOps.get(0).toString()); Assert.assertEquals("SET mykey ", redisOps.get(1).toString()); - byte[][] bitmapKey2 = redisOps.get(2).buildRawOpArgs(); + byte[][] bitmapKey2 = redisOps.get(3).buildRawOpArgs(); Assert.assertEquals("bitmap_key2", new String(bitmapKey2[1])); Assert.assertEquals(700 / 8 + (700 % 8 == 0 ? 0 : 1),bitmapKey2[2].length); - byte[][] bitmapKey = redisOps.get(3).buildRawOpArgs(); + byte[][] bitmapKey = redisOps.get(5).buildRawOpArgs(); Assert.assertEquals("bitmap_key", new String(bitmapKey[1])); Assert.assertEquals(700 / 8 + (700 % 8 == 0 ? 0 : 1),bitmapKey[2].length); Assert.assertNotEquals(0, bitmapKey[2][87] & (1 << 3)); // getbit bitmap_key 700 -> 1 Assert.assertNotEquals(0, bitmapKey[2][8] & (1 << 1)); // getbit bitmap_key 70 -> 1 - Assert.assertEquals("SET common_key value", redisOps.get(4).toString()); + Assert.assertEquals("SET common_key value", redisOps.get(7).toString()); } @Test From 758d1ad5b095df1ad42aa5a72d2247fe63e6376f Mon Sep 17 00:00:00 2001 From: TB Date: Mon, 20 Jul 2026 19:13:06 +0800 Subject: [PATCH 4/7] support 8.2 new command --- .../core/redis/operation/RedisOpType.java | 16 +- .../parser/RedisOpMultiKeysEnum.java | 4 + .../parser/RedisOpMultiKeysParser.java | 20 +- .../parser/RedisOpWithSubKeysEnum.java | 9 +- .../parser/RedisOpWithSubKeysParser.java | 86 ++++- .../parser/GeneralRedisOpParserTest.java | 301 +++++++++++++++++- .../src/test/resources/log4j2.xml | 1 + redis/redis-keeper/dump.rdb | Bin 92 -> 0 bytes 8 files changed, 412 insertions(+), 25 deletions(-) delete mode 100644 redis/redis-keeper/dump.rdb diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/RedisOpType.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/RedisOpType.java index 20e42846c7..e4b54a6cc4 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/RedisOpType.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/RedisOpType.java @@ -40,6 +40,8 @@ public enum RedisOpType { LREM(false, 4), LSET(false, 4), LTRIM(false, 4), + BLMPOP(true,-5), + LMPOP(true,-4), // Hash single HDEL(false, -3), @@ -47,8 +49,15 @@ public enum RedisOpType { HINCRBYFLOAT(false, 4), HMSET(false, -4), HSET(false, -4), - HSETEX(false, -4), + HSETEX(false, -6), HSETNX(false, 4), + HEXPIREAT(false,-6), + HPEXPIREAT(false,-6), + HEXPIRE(false,-6), + HPEXPIRE(false,-6), + HGETDEL(false,-5), + HPERSIST(false,-5), + HGETEX(false,-5), // Set single SADD(false, -3), @@ -62,6 +71,8 @@ public enum RedisOpType { ZREMRANGEBYLEX(false, 4), ZREMRANGEBYRANK(false, 4), ZREMRANGEBYSCORE(false, 4), + ZMPOP(true,-4), + BZMPOP(true,-5), // Stream single XADD(false, -5), @@ -69,6 +80,9 @@ public enum RedisOpType { XSETID(false, 3), XGROUP(false, -2), XCLAIM(false, -6), + XDELEX(false,-6), + XACKDEL(false,-7), + // TTL single EXPIRE(false, 3), diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpMultiKeysEnum.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpMultiKeysEnum.java index 0d501d1e0e..e2a822cd0d 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpMultiKeysEnum.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpMultiKeysEnum.java @@ -13,6 +13,10 @@ public enum RedisOpMultiKeysEnum { MSETNX(RedisOpType.MSETNX, 1, 2), DEL(RedisOpType.DEL, 1, 1), UNLINK(RedisOpType.UNLINK, 1, 1), + LMPOP(RedisOpType.LMPOP, 2, 1), + BLMPOP(RedisOpType.BLMPOP, 3, 1), + ZMPOP(RedisOpType.ZMPOP, 2, 1), + BZMPOP(RedisOpType.BZMPOP, 3, 1), //crdt, CRDT_DEL_REG(RedisOpType.CRDT_DEL_REG, 1, 1), diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpMultiKeysParser.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpMultiKeysParser.java index 6f062e78e8..7961a9f280 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpMultiKeysParser.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpMultiKeysParser.java @@ -6,8 +6,10 @@ import com.ctrip.xpipe.redis.core.redis.operation.RedisOpType; import com.ctrip.xpipe.redis.core.redis.operation.op.RedisOpMultiKVs; import com.ctrip.xpipe.tuple.Pair; +import org.checkerframework.checker.units.qual.A; import java.util.ArrayList; +import java.util.Arrays; import java.util.List; /** @@ -17,6 +19,12 @@ */ public class RedisOpMultiKeysParser extends AbstractRedisOpParser implements RedisOpParser { + private static final byte[] MIN_BYTES = "MIN".getBytes(); + private static final byte[] MIN_BYTES_LOWER = "min".getBytes(); + + private static final byte[] MAX_BYTES = "MAX".getBytes(); + private static final byte[] MAX_BYTES_LOWER = "max".getBytes(); + private RedisOpType redisOpType; private Integer keyStartIndex; private Integer kvNum; @@ -37,7 +45,17 @@ public RedisOp parse(byte[][] args) { List> kvs = new ArrayList<>(); int i = keyStartIndex; - while (i < args.length) { + int argLen = args.length; + if(keyStartIndex>1){ + int numKeys; + try { + numKeys = Integer.parseInt(new String(args[keyStartIndex-1])); + } catch (NumberFormatException e) { + throw new IllegalArgumentException("Invalid numkeys for " + redisOpType); + } + argLen = i+numKeys; + } + while (i < argLen) { RedisKey key = new RedisKey(args[i++]); kvs.add(kvNum == 1 ? Pair.of(key, null) : Pair.of(key, args[i++])); } diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpWithSubKeysEnum.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpWithSubKeysEnum.java index b62e9dc2fb..58831bcc9b 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpWithSubKeysEnum.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpWithSubKeysEnum.java @@ -15,7 +15,14 @@ public enum RedisOpWithSubKeysEnum { HINCRBYFLOAT(RedisOpType.HINCRBYFLOAT, 1, 2,false), HMSET(RedisOpType.HMSET, 1, 2,false), HSET(RedisOpType.HSET, 1, 2,false), - HSETEX(RedisOpType.HSET, 1, 2,false), + HSETEX(RedisOpType.HSETEX, 1, 2,false), + HGETEX(RedisOpType.HGETEX, 1, 1,false), + HGETDEL(RedisOpType.HGETDEL, 1, 1,false), + HEXPIRE(RedisOpType.HEXPIRE, 1, 1,false), + HEXPIREAT(RedisOpType.HEXPIREAT, 1, 1,false), + HPEXPIRE(RedisOpType.HPEXPIRE, 1, 1,false), + HPEXPIREAT(RedisOpType.HPEXPIREAT, 1, 1,false), + HPERSIST(RedisOpType.HPERSIST, 1, 1,false), HSETNX(RedisOpType.HSETNX, 1, 2,false), ZADD(RedisOpType.ZADD, 1, 2,true), diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpWithSubKeysParser.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpWithSubKeysParser.java index 1c649787c5..bdb4887078 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpWithSubKeysParser.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpWithSubKeysParser.java @@ -4,13 +4,11 @@ import com.ctrip.xpipe.redis.core.redis.operation.RedisOp; import com.ctrip.xpipe.redis.core.redis.operation.RedisOpParser; import com.ctrip.xpipe.redis.core.redis.operation.RedisOpType; -import com.ctrip.xpipe.redis.core.redis.operation.op.RedisOpMultiKVs; import com.ctrip.xpipe.redis.core.redis.operation.op.RedisOpMultiSubKey; import com.ctrip.xpipe.tuple.Pair; -import java.util.ArrayList; -import java.util.Arrays; -import java.util.List; +import java.nio.ByteBuffer; +import java.util.*; /** * @author tb @@ -24,12 +22,28 @@ public class RedisOpWithSubKeysParser extends AbstractRedisOpParser implements R private Integer kvNum; private boolean kvReverse; - private byte[] NX_BYTES = new byte[]{'N','X'}; - private byte[] XX_BYTES = new byte[]{'X','X'}; - private byte[] GT_BYTES = new byte[]{'G','T'}; - private byte[] LT_BYTES = new byte[]{'L','T'}; - private byte[] CH_BYTES = new byte[]{'C','H'}; - private byte[] INCR_BYTES = new byte[]{'I','N','C','R'}; + private static final Set NON_KEY_COMMANDS = new HashSet<>(); + private static final Set NON_KEY_COMMANDS_WITH_ARGS = new HashSet<>(); + + private static final byte[] FIELDS_BYTES = "FIELDS".getBytes(); + private static final byte[] FIELDS_BYTES_LOWER = "fields".getBytes(); + + static { + // 无参数关键字 + String[] noArgCmds = { + "NX", "nx", "XX", "xx", "GT", "gt", "LT", "lt", "CH", "ch", + "INCR", "incr", "KEEPTTL", "keepttl", "FNX", "fnx", "FXX", "fxx", "PERSIST", "persist" + }; + for (String cmd : noArgCmds) { + NON_KEY_COMMANDS.add(ByteBuffer.wrap(cmd.getBytes())); + } + // 带一个参数的关键字 + String[] oneArgCmds = { "EX", "ex", "PX", "px", "EXAT", "exat", "PXAT", "pxat" }; + for (String cmd : oneArgCmds) { + NON_KEY_COMMANDS_WITH_ARGS.add(ByteBuffer.wrap(cmd.getBytes())); + NON_KEY_COMMANDS.add(ByteBuffer.wrap(cmd.getBytes())); // 也在无参数集合中,因为需要先识别 + } + } public RedisOpWithSubKeysParser(RedisOpType redisOpType, int keyStartIndex, int kvNum,boolean kvReverse) { this.redisOpType = redisOpType; @@ -43,6 +57,13 @@ public RedisOp parse(byte[][] args) { Pair pair = redisOpType.transfer(redisOpType, args); args = pair.getValue(); + if(usesFieldsSyntax(redisOpType)) { + return parseFieldsCommand(redisOpType,args); + } + return parseGeneralCommands(args, pair); + } + + private RedisOpMultiSubKey parseGeneralCommands(byte[][] args, Pair pair) { // 计算容量 - 主键 + 子键数量 int subKeyCount = (args.length - keyStartIndex-1) / kvNum; int capacity = 1 + subKeyCount; @@ -76,7 +97,7 @@ public RedisOp parse(byte[][] args) { subKeys.add(subKey); } - return new RedisOpMultiSubKey(pair.getKey(), args, key,subKeys); + return new RedisOpMultiSubKey(pair.getKey(), args, key, subKeys); } @Override @@ -85,8 +106,45 @@ public int getOrder() { } private boolean nonKey(byte[] args){ - return Arrays.equals(NX_BYTES,args) || Arrays.equals(XX_BYTES,args) - || Arrays.equals(GT_BYTES,args) || Arrays.equals(LT_BYTES,args) - || Arrays.equals(CH_BYTES,args) || Arrays.equals(INCR_BYTES,args); + return NON_KEY_COMMANDS.contains(ByteBuffer.wrap(args)); + } + + private boolean usesFieldsSyntax(RedisOpType opType) { + return opType == RedisOpType.HSETEX || opType == RedisOpType.HEXPIRE + || opType == RedisOpType.HEXPIREAT || opType == RedisOpType.HGETEX + || opType == RedisOpType.HGETDEL || opType == RedisOpType.HPEXPIRE + || opType == RedisOpType.HPEXPIREAT || opType == RedisOpType.HPERSIST; + } + + private RedisOp parseFieldsCommand(RedisOpType opType, byte[][] args) { + int idx = keyStartIndex; + RedisKey key = new RedisKey(args[idx++]); + + // 跳过 FIELDS 之前的所有可选标志 + while (idx < args.length) { + ByteBuffer tokenBuf = ByteBuffer.wrap(args[idx]); + if (NON_KEY_COMMANDS.contains(tokenBuf)) { + idx++; + if (NON_KEY_COMMANDS_WITH_ARGS.contains(tokenBuf) && idx < args.length) { + idx++; // 跳过标志的参数 + } + } else if (Arrays.equals(args[idx], FIELDS_BYTES) || Arrays.equals(args[idx], FIELDS_BYTES_LOWER)) { + idx++; + break; + } else { + idx++; + } + } + + if (idx >= args.length) throw new IllegalArgumentException("Missing numfields"); + int numFields = Integer.parseInt(new String(args[idx++])); + List subKeys = new ArrayList<>(numFields); + + for (int i = 0; i < numFields; i++) { + if (idx >= args.length) throw new IllegalArgumentException("Incomplete field list"); + subKeys.add(new RedisKey(args[idx])); + idx += kvNum; + } + return new RedisOpMultiSubKey(opType, args, key, subKeys); } } diff --git a/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/parser/GeneralRedisOpParserTest.java b/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/parser/GeneralRedisOpParserTest.java index 69520e8d9f..8aa08a1e39 100644 --- a/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/parser/GeneralRedisOpParserTest.java +++ b/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/parser/GeneralRedisOpParserTest.java @@ -9,9 +9,7 @@ import org.junit.runner.RunWith; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; -import java.util.ArrayList; -import java.util.Arrays; -import java.util.List; +import java.util.*; /** * @author lishanglin @@ -192,11 +190,21 @@ public void testParseAllCmds() { "hincrbyfloat", "hmset", "hset", + "hsetex", "hsetnx", + "hexpireat", + "hpexpireat", + "hpexpire", + "hexpire", + "hgetdel", + "HPERSIST", + "HGETEX", "incr", "incrby", "linsert", "lpop", + "BLMPOP", + "LMPOP", "lpush", "lpushx", "lrem", @@ -222,6 +230,8 @@ public void testParseAllCmds() { "srem", "unlink", "zadd", + "ZMPOP", + "BZMPOP", "zincrby", "zrem", "zremrangebylex", @@ -235,16 +245,47 @@ public void testParseAllCmds() { "script", "multi" ); + // 所有使用 FIELDS 语法的命令(需要特殊构造参数) + Set fieldsCmds = new HashSet<>(Arrays.asList( + "hsetex", "hgetex", "hgetdel", "hpersist" + )); + + Set fieldsCmdWithArgs = new HashSet<>(Arrays.asList( + "hexpire", "hexpireat", + "hpexpire", "hpexpireat" + )); for (String cmd : cmdList) { RedisOpType redisOpType = RedisOpType.lookup(cmd); List parserList = new ArrayList<>(); - System.out.println(cmd); - parserList.add(cmd); - for (int i = 0; i < Math.abs(redisOpType.getArity()) - 1; i++) { - parserList.add("0"); + parserList.add(cmd); // 第一个元素是命令本身 + + if (fieldsCmds.contains(cmd.toLowerCase())) { + // 构造符合 FIELDS 语法的合法参数序列 + parserList.add("mykey"); // 主键 + parserList.add("FIELDS"); + parserList.add("1"); // numfields = 1 + parserList.add("field1"); // 一个 field + if ("hsetex".equalsIgnoreCase(cmd)) { + parserList.add("value1"); // HSETEX 需要 field-value 对 + } + } else if(fieldsCmdWithArgs.contains(cmd.toLowerCase())) { + parserList.add("mykey"); + parserList.add("100"); + parserList.add("FIELDS"); + parserList.add("1"); // numfields = 1 + parserList.add("field1"); + } else { + // 原有逻辑:根据 arity 填充占位参数 + int arity = redisOpType.getArity(); + int minArgs = Math.abs(arity) - 1; // 减去命令名自身 + for (int i = 0; i < minArgs; i++) { + parserList.add("0"); + } } - parser.parse(parserList.toArray()); + + System.out.println(cmd); + parser.parse(parserList.toArray()); // 解析不应抛出异常 } } @@ -260,6 +301,13 @@ public void testGTIDHashCmdWithSubKeys(){ Assert.assertEquals("a1:10", redisOp.getOpGtid()); } + @Test + public void testGTIDHashExCmdWithSubKeys(){ + RedisOp redisOp = parser.parse(Arrays.asList("GTID", "a1:10", "0", "HSETEX","hash1","nx","ex","3600","fields","2", "k1", "v1", "k2", "v2").toArray()); + Assert.assertEquals(RedisOpType.HSETEX, redisOp.getOpType()); + Assert.assertEquals("a1:10", redisOp.getOpGtid()); + } + @Test public void testGTIDSetCmdWithSubKeys(){ RedisOp redisOp = parser.parse(Arrays.asList("GTID", "a1:10", "0", "SADD","set1", "m1","m2").toArray()); @@ -282,6 +330,15 @@ public void testGTIDZSetAddCmdWithSubKeys(){ Assert.assertEquals("a1:10", redisOp.getOpGtid()); } + @Test + public void testGTIDZSetAddCmdWithSubKeysWithXxargs(){ + RedisOp redisOp = parser.parse(Arrays.asList("GTID", "a1:10", "0", "ZADD","zset1","xx","ch","INCR", "1000","zm1","2000","zm2").toArray()); + + Assert.assertEquals(RedisOpType.ZADD, redisOp.getOpType()); + Assert.assertEquals("a1:10", redisOp.getOpGtid()); + } + + @Test public void testGTIDZGeoCmdWithSubKeys(){ @@ -301,4 +358,232 @@ public void testGTIDTransactionCmdWithSubKeys(){ RedisOp execRedisOp = parser.parse(Arrays.asList("GTID", "a1:10", "0", "exec").toArray()); Assert.assertEquals("a1:10", execRedisOp.getOpGtid()); } + + // ---------- 其他 FIELDS 型 Hash 命令测试 ---------- + + @Test + public void testGTIDHexPireCmds() { + // HEXPIRE key seconds FIELDS numfields field + RedisOp redisOp = parser.parse( + Arrays.asList("GTID", "a1:10", "0", "HEXPIRE", "myhash", "100", "FIELDS", "1", "f1").toArray() + ); + Assert.assertEquals(RedisOpType.HEXPIRE, redisOp.getOpType()); + Assert.assertEquals("a1:10", redisOp.getOpGtid()); + + // HEXPIREAT key unix-time-seconds FIELDS numfields field + redisOp = parser.parse( + Arrays.asList("GTID", "a1:10", "0", "HEXPIREAT", "myhash", "123456789", "FIELDS", "2", "f1", "f2").toArray() + ); + Assert.assertEquals(RedisOpType.HEXPIREAT, redisOp.getOpType()); + Assert.assertEquals("a1:10", redisOp.getOpGtid()); + + // HPEXPIRE key milliseconds FIELDS numfields field + redisOp = parser.parse( + Arrays.asList("GTID", "a1:10", "0", "HPEXPIRE", "myhash", "5000", "FIELDS", "1", "f1").toArray() + ); + Assert.assertEquals(RedisOpType.HPEXPIRE, redisOp.getOpType()); + Assert.assertEquals("a1:10", redisOp.getOpGtid()); + + // HPEXPIREAT key unix-time-milliseconds FIELDS numfields field + redisOp = parser.parse( + Arrays.asList("GTID", "a1:10", "0", "HPEXPIREAT", "myhash", "123456789000", "FIELDS", "1", "f1").toArray() + ); + Assert.assertEquals(RedisOpType.HPEXPIREAT, redisOp.getOpType()); + Assert.assertEquals("a1:10", redisOp.getOpGtid()); + } + + @Test + public void testGTIDHgetDelAndHgetex() { + // HGETDEL key FIELDS numfields field + RedisOp redisOp = parser.parse( + Arrays.asList("GTID", "a1:10", "0", "HGETDEL", "myhash", "FIELDS", "1", "f1").toArray() + ); + Assert.assertEquals(RedisOpType.HGETDEL, redisOp.getOpType()); + Assert.assertEquals("a1:10", redisOp.getOpGtid()); + + // HGETEX key PERSIST FIELDS numfields field + redisOp = parser.parse( + Arrays.asList("GTID", "a1:10", "0", "HGETEX", "myhash", "PERSIST", "FIELDS", "1", "f1").toArray() + ); + Assert.assertEquals(RedisOpType.HGETEX, redisOp.getOpType()); + Assert.assertEquals("a1:10", redisOp.getOpGtid()); + + // HGETEX key EX 10 FIELDS 2 field1 field2 + redisOp = parser.parse( + Arrays.asList("GTID", "a1:10", "0", "HGETEX", "myhash", "EX", "10", "FIELDS", "2", "f1", "f2").toArray() + ); + Assert.assertEquals(RedisOpType.HGETEX, redisOp.getOpType()); + Assert.assertEquals("a1:10", redisOp.getOpGtid()); + } + + @Test + public void testGTIDHPersist() { + // HPERSIST key FIELDS numfields field + RedisOp redisOp = parser.parse( + Arrays.asList("GTID", "a1:10", "0", "HPERSIST", "myhash", "FIELDS", "2", "f1", "f2").toArray() + ); + Assert.assertEquals(RedisOpType.HPERSIST, redisOp.getOpType()); + Assert.assertEquals("a1:10", redisOp.getOpGtid()); + } + + @Test + public void testParseOtherHashFieldsCmds() { + // HEXPIRE + + RedisOp redisOp = parser.parse( + Arrays.asList("HEXPIRE", "myhash", "100", "FIELDS", "1", "f1").toArray() + ); + Assert.assertEquals(RedisOpType.HEXPIRE, redisOp.getOpType()); + Assert.assertNull(redisOp.getOpGtid()); + + redisOp = parser.parse( + Arrays.asList("HEXPIRE", "myhash", "100", "xx","FIELDS", "1", "f1").toArray() + ); + Assert.assertEquals(RedisOpType.HEXPIRE, redisOp.getOpType()); + Assert.assertNull(redisOp.getOpGtid()); + + redisOp = parser.parse( + Arrays.asList("HEXPIRE", "myhash", "100", "XX","FIELDS", "1", "f1").toArray() + ); + Assert.assertEquals(RedisOpType.HEXPIRE, redisOp.getOpType()); + Assert.assertNull(redisOp.getOpGtid()); + + // HGETDEL + redisOp = parser.parse( + Arrays.asList("HGETDEL", "myhash", "FIELDS", "2", "f1", "f2").toArray() + ); + Assert.assertEquals(RedisOpType.HGETDEL, redisOp.getOpType()); + Assert.assertNull(redisOp.getOpGtid()); + + // HPERSIST + redisOp = parser.parse( + Arrays.asList("HPERSIST", "myhash", "FIELDS", "1", "f1").toArray() + ); + Assert.assertEquals(RedisOpType.HPERSIST, redisOp.getOpType()); + Assert.assertNull(redisOp.getOpGtid()); + } + + @Test + public void testZMPOPParse() { + RedisOp redisOp = parser.parse( + Arrays.asList("ZMPOP", "2", "zset1", "zset2", "MIN", "COUNT", "5").toArray() + ); + Assert.assertEquals(RedisOpType.ZMPOP, redisOp.getOpType()); + Assert.assertNull(redisOp.getOpGtid()); + + RedisMultiKeyOp multiKeyOp = (RedisMultiKeyOp) redisOp; + Assert.assertEquals(2, multiKeyOp.getKeys().size()); + Assert.assertEquals(new RedisKey("zset1"), multiKeyOp.getKeyValue(0).getKey()); + Assert.assertNull(multiKeyOp.getKeyValue(0).getValue()); + Assert.assertEquals(new RedisKey("zset2"), multiKeyOp.getKeyValue(1).getKey()); + Assert.assertNull(multiKeyOp.getKeyValue(1).getValue()); + } + + @Test + public void testZMPOPWithOneKey() { + RedisOp redisOp = parser.parse( + Arrays.asList("ZMPOP", "1", "zset1", "MAX").toArray() + ); + Assert.assertEquals(RedisOpType.ZMPOP, redisOp.getOpType()); + + RedisMultiKeyOp multiKeyOp = (RedisMultiKeyOp) redisOp; + Assert.assertEquals(1, multiKeyOp.getKeys().size()); + Assert.assertEquals(new RedisKey("zset1"), multiKeyOp.getKeyValue(0).getKey()); + } + + @Test + public void testBZMPOPParse() { + RedisOp redisOp = parser.parse( + Arrays.asList("BZMPOP", "0.5", "2", "zset1", "zset2", "MIN", "COUNT", "3").toArray() + ); + Assert.assertEquals(RedisOpType.BZMPOP, redisOp.getOpType()); + Assert.assertNull(redisOp.getOpGtid()); + + RedisMultiKeyOp multiKeyOp = (RedisMultiKeyOp) redisOp; + Assert.assertEquals(2, multiKeyOp.getKeys().size()); + Assert.assertEquals(new RedisKey("zset1"), multiKeyOp.getKeyValue(0).getKey()); + Assert.assertEquals(new RedisKey("zset2"), multiKeyOp.getKeyValue(1).getKey()); + } + + @Test + public void testLMPOPParse() { + RedisOp redisOp = parser.parse( + Arrays.asList("LMPOP", "2", "list1", "list2", "LEFT", "COUNT", "10").toArray() + ); + Assert.assertEquals(RedisOpType.LMPOP, redisOp.getOpType()); + + RedisMultiKeyOp multiKeyOp = (RedisMultiKeyOp) redisOp; + Assert.assertEquals(2, multiKeyOp.getKeys().size()); + Assert.assertEquals(new RedisKey("list1"), multiKeyOp.getKeyValue(0).getKey()); + Assert.assertEquals(new RedisKey("list2"), multiKeyOp.getKeyValue(1).getKey()); + } + + @Test + public void testBLMPOPParse() { + RedisOp redisOp = parser.parse( + Arrays.asList("BLMPOP", "0.1", "2", "list1", "list2", "RIGHT").toArray() + ); + Assert.assertEquals(RedisOpType.BLMPOP, redisOp.getOpType()); + + RedisMultiKeyOp multiKeyOp = (RedisMultiKeyOp) redisOp; + Assert.assertEquals(2, multiKeyOp.getKeys().size()); + Assert.assertEquals(new RedisKey("list1"), multiKeyOp.getKeyValue(0).getKey()); + Assert.assertEquals(new RedisKey("list2"), multiKeyOp.getKeyValue(1).getKey()); + } + + // ---------- 带 GTID 的测试 ---------- + @Test + public void testGTIDZMPOPParse() { + RedisOp redisOp = parser.parse( + Arrays.asList("GTID", "a1:10", "0", "ZMPOP", "2", "zset1", "zset2", "MIN").toArray() + ); + Assert.assertEquals(RedisOpType.ZMPOP, redisOp.getOpType()); + Assert.assertEquals("a1:10", redisOp.getOpGtid()); + + RedisMultiKeyOp multiKeyOp = (RedisMultiKeyOp) redisOp; + Assert.assertEquals(2, multiKeyOp.getKeys().size()); + Assert.assertEquals(new RedisKey("zset1"), multiKeyOp.getKeyValue(0).getKey()); + } + + @Test + public void testGTIDBZMPOPParse() { + RedisOp redisOp = parser.parse( + Arrays.asList("GTID", "a1:10", "0", "BZMPOP", "0.5", "1", "zset1", "MAX", "COUNT", "5").toArray() + ); + Assert.assertEquals(RedisOpType.BZMPOP, redisOp.getOpType()); + Assert.assertEquals("a1:10", redisOp.getOpGtid()); + + RedisMultiKeyOp multiKeyOp = (RedisMultiKeyOp) redisOp; + Assert.assertEquals(1, multiKeyOp.getKeys().size()); + Assert.assertEquals(new RedisKey("zset1"), multiKeyOp.getKeyValue(0).getKey()); + } + + @Test + public void testGTIDLMPOPParse() { + RedisOp redisOp = parser.parse( + Arrays.asList("GTID", "a1:10", "0", "LMPOP", "3", "l1", "l2", "l3", "RIGHT", "COUNT", "2").toArray() + ); + Assert.assertEquals(RedisOpType.LMPOP, redisOp.getOpType()); + Assert.assertEquals("a1:10", redisOp.getOpGtid()); + + RedisMultiKeyOp multiKeyOp = (RedisMultiKeyOp) redisOp; + Assert.assertEquals(3, multiKeyOp.getKeys().size()); + Assert.assertEquals(new RedisKey("l1"), multiKeyOp.getKeyValue(0).getKey()); + Assert.assertEquals(new RedisKey("l2"), multiKeyOp.getKeyValue(1).getKey()); + Assert.assertEquals(new RedisKey("l3"), multiKeyOp.getKeyValue(2).getKey()); + } + + @Test + public void testGTIDBLMPOPParse() { + RedisOp redisOp = parser.parse( + Arrays.asList("GTID", "a1:10", "0", "BLMPOP", "2", "2", "l1", "l2", "LEFT").toArray() + ); + Assert.assertEquals(RedisOpType.BLMPOP, redisOp.getOpType()); + Assert.assertEquals("a1:10", redisOp.getOpGtid()); + + RedisMultiKeyOp multiKeyOp = (RedisMultiKeyOp) redisOp; + Assert.assertEquals(2, multiKeyOp.getKeys().size()); + Assert.assertEquals(new RedisKey("l1"), multiKeyOp.getKeyValue(0).getKey()); + Assert.assertEquals(new RedisKey("l2"), multiKeyOp.getKeyValue(1).getKey()); + } } diff --git a/redis/redis-integration-test/src/test/resources/log4j2.xml b/redis/redis-integration-test/src/test/resources/log4j2.xml index 0a2be948bb..43806144cf 100644 --- a/redis/redis-integration-test/src/test/resources/log4j2.xml +++ b/redis/redis-integration-test/src/test/resources/log4j2.xml @@ -23,6 +23,7 @@ + diff --git a/redis/redis-keeper/dump.rdb b/redis/redis-keeper/dump.rdb deleted file mode 100644 index 99c56e7b9897ef94120256f032024a21b87e3914..0000000000000000000000000000000000000000 GIT binary patch literal 0 HcmV?d00001 literal 92 zcmWG?b@2=~Ffg$E#aWb^l3A= Date: Tue, 21 Jul 2026 14:03:07 +0800 Subject: [PATCH 5/7] support multi split --- .../redis/operation/op/RedisOpItemParser.java | 2 +- .../redis/rdb/parser/RdbQuickList2Parser.java | 2 +- .../operation/op/RedisOpItemParserTest.java | 39 +++++++ ...erServerToKeeperToFakeXsyncServerTest.java | 98 +++++++++++++++++- .../src/test/resources/applier/dump6379.rdb | Bin 0 -> 482 bytes .../applier/command/DefaultRawCommand.java | 57 ++++++++++ .../command/TransactionAsyncCommand.java | 87 +++++++--------- .../sequence/DefaultSequenceController.java | 50 ++++----- .../sync/DefaultCommandDispatcher.java | 3 +- 9 files changed, 255 insertions(+), 83 deletions(-) create mode 100644 redis/redis-integration-test/src/test/resources/applier/dump6379.rdb create mode 100644 redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/command/DefaultRawCommand.java diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/op/RedisOpItemParser.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/op/RedisOpItemParser.java index 5d6cc08086..4bee46f6a7 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/op/RedisOpItemParser.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/op/RedisOpItemParser.java @@ -16,7 +16,7 @@ private RedisOpItemParser() { public static RedisOpItem parse(RedisOpParser redisOpParser, Object[] payload) { RedisOpItem redisOpItem = new RedisOpItem(); try { - RedisOp redisOp = redisOpParser.parse(payload); + RedisOp redisOp = redisOpParser.parse(payload); redisOpItem.setRedisOpType(redisOp.getOpType()); redisOpItem.setGtid(redisOp.getOpGtid()); redisOpItem.setDbId(redisOp.getDbId()); diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/RdbQuickList2Parser.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/RdbQuickList2Parser.java index 2978939f59..360d8b4b14 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/RdbQuickList2Parser.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/parser/RdbQuickList2Parser.java @@ -63,8 +63,8 @@ public Integer read(ByteBuf byteBuf) { break; case READ_DATA: byte[] data = rdbStringParser.read(byteBuf); - rdbStringParser.reset(); if (data != null) { + rdbStringParser.reset(); readCnt++; if (readCnt >= len.getLenValue()) { diff --git a/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/operation/op/RedisOpItemParserTest.java b/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/operation/op/RedisOpItemParserTest.java index ed351bf893..a8f726d9f9 100644 --- a/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/operation/op/RedisOpItemParserTest.java +++ b/redis/redis-core/src/test/java/com/ctrip/xpipe/redis/core/redis/operation/op/RedisOpItemParserTest.java @@ -29,4 +29,43 @@ public void testParseFromDirectByteBufInOutPayload() throws IOException { Assert.assertNotNull(item.getRedisKey()); Assert.assertArrayEquals("foo".getBytes(), item.getRedisKey().get()); } + + + @Test + public void testParseFromDirectByteBufInOutPayloadZAdd() throws IOException { + DirectByteBufInOutPayload gtid = new DirectByteBufInOutPayload(); + gtid.startInput(); + gtid.in(Unpooled.wrappedBuffer("GTID".getBytes())); + + DirectByteBufInOutPayload gtidStr = new DirectByteBufInOutPayload(); + gtidStr.startInput(); + gtidStr.in(Unpooled.wrappedBuffer("573c4305037d8c0cddb1ee36b5455d3a9d36f43e:1268595202".getBytes())); + + DirectByteBufInOutPayload db = new DirectByteBufInOutPayload(); + db.startInput(); + db.in(Unpooled.wrappedBuffer("0".getBytes())); + + + DirectByteBufInOutPayload command = new DirectByteBufInOutPayload(); + command.startInput(); + command.in(Unpooled.wrappedBuffer("zadd".getBytes())); + DirectByteBufInOutPayload key = new DirectByteBufInOutPayload(); + key.startInput(); + key.in(Unpooled.wrappedBuffer("zset:did2vid:1027cff212d9d9a6".getBytes())); + DirectByteBufInOutPayload nx = new DirectByteBufInOutPayload(); + nx.startInput(); + nx.in(Unpooled.wrappedBuffer("nx".getBytes())); + DirectByteBufInOutPayload score = new DirectByteBufInOutPayload(); + score.startInput(); + score.in(Unpooled.wrappedBuffer("1.784010284E9".getBytes())); + DirectByteBufInOutPayload value = new DirectByteBufInOutPayload(); + value.startInput(); + value.in(Unpooled.wrappedBuffer("09C74910607911F083DBFFC139EB4BA5".getBytes())); + + RedisOpItem item = RedisOpItemParser.parse(parser, new Object[]{gtid,gtidStr,db,command, key,nx,score, value}); + + Assert.assertEquals(RedisOpType.ZADD, item.getRedisOpType()); + Assert.assertNotNull(item.getRedisKey()); + Assert.assertArrayEquals("zset:did2vid:1027cff212d9d9a6".getBytes(), item.getRedisKey().get()); + } } diff --git a/redis/redis-integration-test/src/test/java/com/ctrip/xpipe/redis/integratedtest/applier/ApplierServerToKeeperToFakeXsyncServerTest.java b/redis/redis-integration-test/src/test/java/com/ctrip/xpipe/redis/integratedtest/applier/ApplierServerToKeeperToFakeXsyncServerTest.java index 8367f8a896..7dc750180b 100644 --- a/redis/redis-integration-test/src/test/java/com/ctrip/xpipe/redis/integratedtest/applier/ApplierServerToKeeperToFakeXsyncServerTest.java +++ b/redis/redis-integration-test/src/test/java/com/ctrip/xpipe/redis/integratedtest/applier/ApplierServerToKeeperToFakeXsyncServerTest.java @@ -1,10 +1,8 @@ package com.ctrip.xpipe.redis.integratedtest.applier; -import com.ctrip.xpipe.redis.core.entity.ApplierMeta; import com.ctrip.xpipe.redis.core.entity.KeeperMeta; import com.ctrip.xpipe.redis.core.entity.RedisMeta; -import com.ctrip.xpipe.redis.keeper.applier.DefaultApplierServer; import org.junit.Assert; import org.junit.Test; import redis.clients.jedis.Jedis; @@ -84,4 +82,100 @@ public void testKeeperApplier2Redis() throws Exception { applier.stop(); } + + @Test + public void testMultiCommandKeeperApplier2Redis() throws Exception { + waitConditionUntilTimeOut(() -> 1 == server.slaveCount()); + for(KeeperMeta keeperMeta : getDcKeepers(dc, getClusterId(), getShardId())) { + waitConditionUntilTimeOut(() -> getRedisKeeperServer(keeperMeta).getRedisMaster().getMasterState().equals(REDIS_REPL_CONNECTED)); + } + server.propagate("multi"); + server.propagate("exec"); + server.propagate("multi"); + server.propagate("set {tag}t1 vt1"); + server.propagate("set {tag}t2 vt2"); + server.propagate("set {tag}t3 vt3"); + server.propagate("exec"); + server.propagate("multi"); + server.propagate("incr in"); + server.propagate("set k1 v11"); + server.propagate("set k2 v22"); + server.propagate("mset k1 v1 k2 v2 k3 v3 k4 v4"); + server.propagate("hset hk1 hf1 hv1 hf2 hv2 hf3 hv3 hf4 hv4"); + server.propagate("set k2 v23"); + server.propagate("incr out"); + server.propagate("GTID f32e4ba72875a76bfd1c92ea1857c4755ea0d680:1 0 exec"); + + server.propagate("GTID f32e4ba72875a76bfd1c92ea1857c4755ea0d680:2 0 hset h1 f1 v1 f2 v2"); + server.propagate("GTID f32e4ba72875a76bfd1c92ea1857c4755ea0d680:3 0 zadd z1 1 v1 2 v2"); + + + sleep(10000); + List redisMetas = getDcRedises(dc,getClusterId(),getShardId()); + int keySize=0; + int bigListLen = 0; + int bigHashLen = 0; + int bigSetLen = 0; + int bigNormalSet1Len = 0; + int bigNormalListLen = 0; + int bigNormalHashLen = 0; + int bigZsetLen = 0; + int bigNormalSetLen = 0; + int h1Len = 0; + int z1Len = 0; + String inCount =""; + String outCount = ""; + String v2 = ""; + String v1 = ""; + for(RedisMeta redisMeta:redisMetas){ + Jedis jedis = new Jedis(redisMeta.getIp(),redisMeta.getPort()); + Set keys = jedis.keys("*"); + keySize += keys.size(); + bigListLen += jedis.llen("biglist"); + bigHashLen += jedis.hlen("bighash"); + bigSetLen += jedis.scard("bigset"); + bigNormalSet1Len += jedis.scard("bignormalset1"); + bigNormalListLen += jedis.llen("bignormallist"); + bigNormalHashLen += jedis.hlen("bignormalhash"); + bigZsetLen += jedis.zcard("bigzset"); + bigNormalSetLen += jedis.zcard("bignormalset"); + h1Len += jedis.hlen("h1"); + z1Len += jedis.zcard("z1"); + String in = jedis.get("in"); + if(in != null){ + inCount = in; + } + String out = jedis.get("out"); + if(out != null){ + outCount = out; + } + String v11 = jedis.get("k1"); + + if(v11 != null){ + v1 = v11; + } + String v22 = jedis.get("k2"); + if(v22 != null) { + v2 = v22; + } + } + Assert.assertEquals(20,keySize); + Assert.assertEquals(3,bigListLen); + Assert.assertEquals(3,bigHashLen); + Assert.assertEquals(3,bigSetLen); + Assert.assertEquals(3,bigNormalSet1Len); + Assert.assertEquals(3,bigNormalListLen); + Assert.assertEquals(3,bigNormalHashLen); + Assert.assertEquals(3,bigZsetLen); + Assert.assertEquals(3,bigNormalSetLen); + Assert.assertEquals(2,h1Len); + Assert.assertEquals(2,z1Len); + Assert.assertEquals("1",inCount); + Assert.assertEquals("1",outCount); + Assert.assertEquals("v23",v2); + Assert.assertEquals("v1",v1); + + applier.stop(); + } + } diff --git a/redis/redis-integration-test/src/test/resources/applier/dump6379.rdb b/redis/redis-integration-test/src/test/resources/applier/dump6379.rdb new file mode 100644 index 0000000000000000000000000000000000000000..2097dc3c6ff2283e295c06cd9e80006318bd656c GIT binary patch literal 482 zcmcJKu};G<5QeP?i?kJe4i;CA6B4pzC<9vsOE%=#7Ziw_D0Zr514K(<2!4}j`zd=?Q-6s|cQz+OipV9~Gwi+Z+ zpmy7rXUEPPs7s$&7F&27xERw-)wnq_73i8j#3NMqO6gAnWofCmGk!qYtz{LGP%a z4cQg1?;fFJNxDKwihjyD4u8|#!E7ojrK!`5%4|4sCH G{Ok)r#%)gk literal 0 HcmV?d00001 diff --git a/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/command/DefaultRawCommand.java b/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/command/DefaultRawCommand.java new file mode 100644 index 0000000000..185d645b82 --- /dev/null +++ b/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/command/DefaultRawCommand.java @@ -0,0 +1,57 @@ +package com.ctrip.xpipe.redis.keeper.applier.command; + +import com.ctrip.xpipe.client.redis.AsyncRedisClient; +import com.ctrip.xpipe.command.AbstractCommand; +import com.ctrip.xpipe.redis.core.redis.operation.RedisMultiKeyOp; +import com.ctrip.xpipe.redis.core.redis.operation.RedisOp; +import com.ctrip.xpipe.redis.core.redis.operation.RedisSingleKeyOp; + +import java.util.List; + +/** + * @author TB + * @date 2026/7/14 11:18 + */ +public class DefaultRawCommand extends AbstractCommand implements RedisOpDataCommand{ + + final AsyncRedisClient client; + + final Object resource; + + final List rawArgs; + + public DefaultRawCommand(AsyncRedisClient client, Object resource, List rawArgs) { + this.client = client; + this.resource = resource; + this.rawArgs = rawArgs; + } + + + @Override + protected void doExecute() throws Throwable { + long startTime = System.nanoTime(); + + client + .writeMulti(resource, 0, rawArgs.toArray()) + .addListener(f -> { + if (getLogger().isDebugEnabled()) { + getLogger().debug("[command] write key {} end, total time {}", redisOp() instanceof RedisSingleKeyOp ? ((RedisSingleKeyOp) redisOp()).getKey() : (redisOp() instanceof RedisMultiKeyOp ? keys() : "none"), System.nanoTime() - startTime); + } + if (f.isSuccess()) { + future().setSuccess(true); + } else { + future().setFailure(f.cause()); + } + }); + } + + @Override + protected void doReset() { + + } + + @Override + public RedisOp redisOp() { + return null; + } +} diff --git a/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/command/TransactionAsyncCommand.java b/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/command/TransactionAsyncCommand.java index 54f5f50521..5adf18ac12 100644 --- a/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/command/TransactionAsyncCommand.java +++ b/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/command/TransactionAsyncCommand.java @@ -1,14 +1,16 @@ package com.ctrip.xpipe.redis.keeper.applier.command; +import com.ctrip.xpipe.api.command.CommandFuture; +import com.ctrip.xpipe.api.command.CommandFutureListener; import com.ctrip.xpipe.client.redis.AsyncRedisClient; +import com.ctrip.xpipe.command.ParallelCommandChain; import com.ctrip.xpipe.redis.core.redis.operation.RedisKey; import com.ctrip.xpipe.redis.core.redis.operation.RedisMultiKeyOp; import com.ctrip.xpipe.redis.core.redis.operation.RedisSingleKeyOp; -import java.util.ArrayList; -import java.util.HashSet; -import java.util.List; -import java.util.Set; +import java.util.*; +import java.util.concurrent.ExecutorService; +import java.util.stream.Collectors; /** * @author TB @@ -24,77 +26,58 @@ public class TransactionAsyncCommand extends TransactionCommand{ private Object resource; - public TransactionAsyncCommand(AsyncRedisClient client){ + final ExecutorService workThreads; + + public TransactionAsyncCommand(AsyncRedisClient client,ExecutorService workThreads){ super(); this.client = client; this.redisKeys = new HashSet<>(); + this.workThreads = workThreads; } public void addTransactionCommands(RedisOpCommand redisOpCommand, long commandOffset, String gtid) { - if(redisOpCommand instanceof MultiDataCommand){ - MultiDataCommand multiDataCommand = (MultiDataCommand) redisOpCommand; - redisKeys.addAll(multiDataCommand.keys()); - }else if(redisOpCommand instanceof DefaultDataCommand){ - DefaultDataCommand defaultDataCommand = (DefaultDataCommand) redisOpCommand; - redisKeys.add(defaultDataCommand.key()); - } super.addTransactionCommands(redisOpCommand,commandOffset,gtid); } @Override protected void doExecute() throws Throwable{ - List multiRawArgs = new ArrayList<>(transactionCommands.size()); - + Map> resourceMultiRawArgs = new HashMap<>(); for(RedisOpCommand redisOpCommand:transactionCommands){ if(redisOpCommand instanceof MultiDataCommand){ MultiDataCommand multiDataCommand = (MultiDataCommand) redisOpCommand; - if(resource == null){ - resource = client.select(multiDataCommand.keys().get(0).get()); - } - if(multiDataCommand.getDbNumber() != 0) { - Object[] selectArgs = new byte[][]{"select".getBytes(), (multiDataCommand.getDbNumber() + "").getBytes()}; - multiRawArgs.add(selectArgs); + List keys = multiDataCommand.keys().stream().map(RedisKey::get).collect(Collectors.toList()); + Map> resourceOps = client.selectMulti(keys); + for(Map.Entry> entry:resourceOps.entrySet()){ + List args = resourceMultiRawArgs.computeIfAbsent(entry.getKey(),(key)-> new ArrayList<>()); + if(multiDataCommand.getDbNumber() != 0) { + byte[][] selectArgs = new byte[][]{"select".getBytes(), (multiDataCommand.getDbNumber() + "").getBytes()}; + args.add(selectArgs); + } + args.add(multiDataCommand.redisOpAsMulti().subOp(entry.getValue().stream().map(keys::indexOf).collect(Collectors.toSet())).buildRawOpArgs()); } - multiRawArgs.add(multiDataCommand.redisOp().buildRawOpArgs()); }else if(redisOpCommand instanceof DefaultDataCommand){ DefaultDataCommand defaultDataCommand = (DefaultDataCommand) redisOpCommand; - if(resource == null) { - resource = client.select(defaultDataCommand.key().get()); - } + Object resource = client.select(defaultDataCommand.key().get()); + List args = resourceMultiRawArgs.computeIfAbsent(resource,(key)-> new ArrayList<>()); if(defaultDataCommand.getDbNumber() != 0) { - Object[] selectArgs = new byte[][]{"select".getBytes(), (defaultDataCommand.getDbNumber() + "").getBytes()}; - multiRawArgs.add(selectArgs); + byte[][] selectArgs = new byte[][]{"select".getBytes(), (defaultDataCommand.getDbNumber() + "").getBytes()}; + args.add(selectArgs); } - multiRawArgs.add(defaultDataCommand.redisOp().buildRawOpArgs()); + args.add(defaultDataCommand.redisOp().buildRawOpArgs()); } } - long startTime = System.nanoTime(); + ParallelCommandChain parallelCommandChain = new ParallelCommandChain(workThreads, false); - client - .writeMulti(resource, 0, multiRawArgs.toArray()) - .addListener(f -> { - if (getLogger().isDebugEnabled()) { - getLogger().debug("[command] write key {} end, total time {}", redisOp() instanceof RedisSingleKeyOp ? ((RedisSingleKeyOp) redisOp()).getKey() : (redisOp() instanceof RedisMultiKeyOp ? keys() : "none"), System.nanoTime() - startTime); - } - if (f.isSuccess()) { - future().setSuccess(true); - } else { - if (f.cause().getMessage().startsWith(ERR_GTID_COMMAND_EXECUTED)) { - future().setSuccess(true); - } else { - future().setFailure(f.cause()); - } - } - }); - } - - public boolean validTransaction(){ - if(redisKeys.size() == 1) return true; - for(RedisKey redisKey:redisKeys) { - byte[] tag = client.hashTag(redisKey.get()); - return tag != null; + for(Map.Entry> entry:resourceMultiRawArgs.entrySet()){ + parallelCommandChain.add(new DefaultRawCommand(client,entry.getKey(),entry.getValue())); } - return false; + parallelCommandChain.execute().addListener(commandFuture -> { + if (commandFuture.isSuccess()) { + TransactionAsyncCommand.this.future().setSuccess(true); + } else { + TransactionAsyncCommand.this.future().setFailure(commandFuture.cause()); + } + }); } } diff --git a/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/sequence/DefaultSequenceController.java b/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/sequence/DefaultSequenceController.java index 2e5e92214b..16f67ab787 100644 --- a/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/sequence/DefaultSequenceController.java +++ b/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/sequence/DefaultSequenceController.java @@ -358,46 +358,32 @@ private void submitObstacle(RedisOpCommand command, long commandOffset, GtidS List transactionOps = redisOpTransactionAdapter.getTransactionOps(); List> dependencies = new ArrayList<>(); Set transactionOpKeys = new HashSet<>(); - RedisKey commandKey = null; + for(RedisOp transactionOp:transactionOps){ if(transactionOp instanceof RedisSingleKeyOp) { RedisSingleKeyOp redisSingleKeyOp = (RedisSingleKeyOp) transactionOp; RedisKey redisKey = redisSingleKeyOp.getKey(); - if(commandKey == null) { - commandKey = redisKey; - } transactionOpKeys.add(redisKey); }else if(transactionOp instanceof RedisMultiKeyOp){ RedisMultiKeyOp redisMultiKeyOp = (RedisMultiKeyOp) transactionOp; for(RedisKey redisKey:redisMultiKeyOp.getKeys()){ - transactionOpKeys.add(redisKey); - if(commandKey == null) { - commandKey = redisKey; + RedisKey commandKey = redisKey; + if (redisKey != null && redisKey.get()[0] == TAG_START) { + commandKey = new RedisKey(client.hashTag(redisKey.get())); + transactionOpKeys.add(commandKey); + break; } + transactionOpKeys.add(commandKey); } }else if(transactionOp instanceof RedisMultiSubKeyOp){ RedisMultiSubKeyOp redisMultiSubKeyOp = (RedisMultiSubKeyOp) transactionOp; RedisKey redisKey = redisMultiSubKeyOp.getKey(); - if(commandKey == null) { - commandKey = redisKey; - } transactionOpKeys.add(redisKey); } } - SequenceCommand lastSameKey = runningCommands.get(commandKey); - if (lastSameKey != null) { - dependencies.add(lastSameKey); - } - if (transactionOpKeys.size() != 1) { - if (commandKey != null && commandKey.get()[0] == TAG_START) { - commandKey = new RedisKey(client.hashTag(commandKey.get())); - } - } - - lastSameKey = multiRunningCommands.get(commandKey); - if (lastSameKey != null) { - dependencies.add(lastSameKey); + for(RedisKey redisKey:transactionOpKeys){ + addIfDependencies(dependencies,redisKey); } /* make command */ @@ -408,9 +394,11 @@ private void submitObstacle(RedisOpCommand command, long commandOffset, GtidS // runningCommands = new HashMap<>(); obstacle = current; - multiRunningCommands.put(commandKey,current); + for(RedisKey commandKey:transactionOpKeys) { + multiRunningCommands.put(commandKey, current); - forgetObstacleWhenDone(current,commandKey); + forgetObstacleWhenDone(current, commandKey); + } /* do some stuff when finish */ @@ -422,6 +410,18 @@ private void submitObstacle(RedisOpCommand command, long commandOffset, GtidS current.execute(); } + private void addIfDependencies(List> dependencies,RedisKey commandKey){ + SequenceCommand lastSameKey = runningCommands.get(commandKey); + if (lastSameKey != null) { + dependencies.add(lastSameKey); + } + + lastSameKey = multiRunningCommands.get(commandKey); + if (lastSameKey != null) { + dependencies.add(lastSameKey); + } + } + private void forgetObstacleWhenDone(SequenceCommand sequenceCommand,RedisKey commandKey) { sequenceCommand.future().addListener((f) -> { if(sequenceCommand == multiRunningCommands.get(commandKey)) { diff --git a/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/sync/DefaultCommandDispatcher.java b/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/sync/DefaultCommandDispatcher.java index 9aa6e7b14f..681d7168a4 100644 --- a/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/sync/DefaultCommandDispatcher.java +++ b/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/sync/DefaultCommandDispatcher.java @@ -177,14 +177,13 @@ protected int toInt(byte[] value) { } private void addTransactionStart(RedisOpCommand multiCommand, long commandOffsetToAccumulate, String gtid) { - transactionCommand.set(new TransactionAsyncCommand(client)); + transactionCommand.set(new TransactionAsyncCommand(client,workerThreads)); transactionCommand.get().addTransactionStart(multiCommand, commandOffsetToAccumulate, gtid); } private void addTransactionEndAndSubmit(RedisOpCommand execCommand, long commandOffsetToAccumulate, String gtid) { transactionCommand.get().addTransactionEnd(execCommand, commandOffsetToAccumulate, gtid); TransactionAsyncCommand command = transactionCommand.getAndSet(null); - if(!command.validTransaction()) throw new RedisRuntimeException("diff keys or no hash tag not supported"); sequenceController.submit(command, command.commandOffset(), command.getGtidSet()); } From 61e39105d7a7dd3a273ee0b15566960bcdda014b Mon Sep 17 00:00:00 2001 From: TB Date: Wed, 22 Jul 2026 15:26:24 +0800 Subject: [PATCH 6/7] RedisOpType add fields attr --- .../core/redis/operation/RedisOpType.java | 26 +++++++++++++------ .../parser/RedisOpWithSubKeysParser.java | 9 +------ .../redis/core/redis/rdb/RdbParseContext.java | 5 +--- 3 files changed, 20 insertions(+), 20 deletions(-) diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/RedisOpType.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/RedisOpType.java index e4b54a6cc4..e69e0e2ded 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/RedisOpType.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/RedisOpType.java @@ -49,15 +49,15 @@ public enum RedisOpType { HINCRBYFLOAT(false, 4), HMSET(false, -4), HSET(false, -4), - HSETEX(false, -6), HSETNX(false, 4), - HEXPIREAT(false,-6), - HPEXPIREAT(false,-6), - HEXPIRE(false,-6), - HPEXPIRE(false,-6), - HGETDEL(false,-5), - HPERSIST(false,-5), - HGETEX(false,-5), + HSETEX(false, true,-6), + HEXPIREAT(false,true,-6), + HPEXPIREAT(false,true,-6), + HEXPIRE(false,true,-6), + HPEXPIRE(false,true,-6), + HGETDEL(false,true,-5), + HPERSIST(false,true,-5), + HGETEX(false,true,-5), // Set single SADD(false, -3), @@ -170,6 +170,8 @@ public enum RedisOpType { private RedisOpCrdtTransfer transfer; + private boolean fields; + private static final Map NAME_CACHE = new HashMap<>(); static { for (RedisOpType op : values()) { @@ -191,6 +193,12 @@ public enum RedisOpType { this(multiKey, arity, swallow, null); } + RedisOpType(boolean multiKey, boolean fields, int artiy) { + this.supportMultiKey = multiKey; + this.fields = fields; + this.arity = artiy; + } + RedisOpType(boolean multiKey, int arity, RedisOpCrdtTransfer transfer) { this(multiKey, arity, false, transfer); } @@ -214,6 +222,8 @@ public boolean isSwallow() { return swallow; } + public boolean isFields() {return fields;} + public Pair transfer(RedisOpType redisOpType, byte[][] args) { if (null == transfer) { return Pair.of(redisOpType, args); diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpWithSubKeysParser.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpWithSubKeysParser.java index bdb4887078..f29ba2eb8d 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpWithSubKeysParser.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpWithSubKeysParser.java @@ -57,7 +57,7 @@ public RedisOp parse(byte[][] args) { Pair pair = redisOpType.transfer(redisOpType, args); args = pair.getValue(); - if(usesFieldsSyntax(redisOpType)) { + if(redisOpType.isFields()) { return parseFieldsCommand(redisOpType,args); } return parseGeneralCommands(args, pair); @@ -109,13 +109,6 @@ private boolean nonKey(byte[] args){ return NON_KEY_COMMANDS.contains(ByteBuffer.wrap(args)); } - private boolean usesFieldsSyntax(RedisOpType opType) { - return opType == RedisOpType.HSETEX || opType == RedisOpType.HEXPIRE - || opType == RedisOpType.HEXPIREAT || opType == RedisOpType.HGETEX - || opType == RedisOpType.HGETDEL || opType == RedisOpType.HPEXPIRE - || opType == RedisOpType.HPEXPIREAT || opType == RedisOpType.HPERSIST; - } - private RedisOp parseFieldsCommand(RedisOpType opType, byte[][] args) { int idx = keyStartIndex; RedisKey key = new RedisKey(args[idx++]); diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbParseContext.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbParseContext.java index 70bea04705..6f09fc92d4 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbParseContext.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/rdb/RdbParseContext.java @@ -144,10 +144,7 @@ enum RdbType { private static final Map> versionTypes = new HashMap<>(); RdbType(short code, boolean rdbOp, Function parserConstructor) { - this.version = 0; - this.code = code; - this.rdbOp = rdbOp; - this.parserConstructor = parserConstructor; + this(0,code,rdbOp,parserConstructor); } RdbType(int version,short code, boolean rdbOp, Function parserConstructor) { From 22abb62b87f60db6932c0f7c8a12a3cc45dbfa05 Mon Sep 17 00:00:00 2001 From: TB Date: Thu, 23 Jul 2026 16:02:17 +0800 Subject: [PATCH 7/7] fix hashtag check & subkey parser merge general fields logic --- .../parser/RedisOpWithSubKeysParser.java | 143 ++++++++++-------- .../sequence/DefaultSequenceController.java | 10 +- 2 files changed, 84 insertions(+), 69 deletions(-) diff --git a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpWithSubKeysParser.java b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpWithSubKeysParser.java index f29ba2eb8d..fa8f0f08c6 100644 --- a/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpWithSubKeysParser.java +++ b/redis/redis-core/src/main/java/com/ctrip/xpipe/redis/core/redis/operation/parser/RedisOpWithSubKeysParser.java @@ -57,47 +57,92 @@ public RedisOp parse(byte[][] args) { Pair pair = redisOpType.transfer(redisOpType, args); args = pair.getValue(); - if(redisOpType.isFields()) { - return parseFieldsCommand(redisOpType,args); + RedisKey key = new RedisKey(args[keyStartIndex]); + int dataStart = keyStartIndex + 1; // 主键之后的数据起始位置 + int subKeyStart; + int subKeyCount; + + if (redisOpType.isFields()) { + // 1. 跳过 FIELDS 之前的所有标志 + int idx = dataStart; + while (idx < args.length) { + if (isFieldsKeyword(args[idx])) { + idx++; + break; + } + int skipped = skipNonKeyTokenIfPresent(args, idx); + if (skipped > idx) { + idx = skipped; + } else { + idx++; + } + } + if (idx >= args.length) throw new IllegalArgumentException("Missing numfields"); + subKeyCount = Integer.parseInt(new String(args[idx++])); + subKeyStart = idx; + } else { + // 2. 普通命令:跳过开头的非键标志(如 ZADD 的 NX/EX 等) + subKeyStart = skipLeadingNonKeyTokens(args, dataStart); + subKeyCount = (args.length - subKeyStart) / kvNum; } - return parseGeneralCommands(args, pair); + + // 3. 共用同一套子键提取逻辑 + List subKeys = extractSubKeys(args, subKeyStart, subKeyCount); + RedisOpType finalOpType = redisOpType.isFields() ? redisOpType : pair.getKey(); + return new RedisOpMultiSubKey(finalOpType, args, key, subKeys); } - private RedisOpMultiSubKey parseGeneralCommands(byte[][] args, Pair pair) { - // 计算容量 - 主键 + 子键数量 - int subKeyCount = (args.length - keyStartIndex-1) / kvNum; - int capacity = 1 + subKeyCount; - List subKeys = new ArrayList<>(capacity); - - int i = keyStartIndex; - RedisKey key = new RedisKey(args[i++]); - // 处理子键 - while (i < args.length) { - RedisKey subKey = null; - - if (!kvReverse) { - // 正常顺序:子键在前 - subKey = new RedisKey(args[i]); - i += kvNum; // 跳过值部分 - } else { + /** + * 统一子键提取:从 start 开始,每次取一个子键,然后按 kvNum 跳过值部分。 + * 支持 kvReverse(反向模式)和 ZADD 非键标志跳过。 + */ + private List extractSubKeys(byte[][] args, int start, int count) { + List subKeys = new ArrayList<>(count); + int i = start; + for (int k = 0; k < count; k++) { + if (kvReverse) { // 反向顺序:值在前,子键在后 - if (kvNum == 2) { - if(RedisOpType.ZADD == redisOpType && nonKey(args[i])){ - i++; + if (kvNum == 2 && RedisOpType.ZADD == redisOpType) { + int skipped = skipNonKeyTokenIfPresent(args, i); + if (skipped > i) { + i = skipped; + k--; // 未消耗子键,重新处理当前子键 continue; } - subKey = new RedisKey(args[i + 1]); // 跳过第一个值,取第二个作为子键 - i += 2; - } else if (kvNum == 3) { - subKey = new RedisKey(args[i + 2]); // 跳过前两个值,取第三个作为子键 - i += 3; } + int keyOffset = kvNum - 1; // kvNum=2 → offset=1, kvNum=3 → offset=2 + subKeys.add(new RedisKey(args[i + keyOffset])); + i += kvNum; + } else { + subKeys.add(new RedisKey(args[i])); + i += kvNum; } + } + return subKeys; + } - subKeys.add(subKey); + /** 跳过从 index 开始的所有连续非键标志,返回第一个非标志的索引 */ + private int skipLeadingNonKeyTokens(byte[][] args, int index) { + while (index < args.length) { + int skipped = skipNonKeyTokenIfPresent(args, index); + if (skipped == index) break; + index = skipped; } + return index; + } - return new RedisOpMultiSubKey(pair.getKey(), args, key, subKeys); + /** 跳过单个非键标志(及其可能的参数),返回新的索引 */ + private int skipNonKeyTokenIfPresent(byte[][] args, int index) { + if (index >= args.length || !nonKey(args[index])) return index; + index++; + if (index < args.length && NON_KEY_COMMANDS_WITH_ARGS.contains(ByteBuffer.wrap(args[index - 1]))) { + index++; + } + return index; + } + + private boolean isFieldsKeyword(byte[] arg) { + return Arrays.equals(arg, FIELDS_BYTES) || Arrays.equals(arg, FIELDS_BYTES_LOWER); } @Override @@ -105,39 +150,7 @@ public int getOrder() { return 0; } - private boolean nonKey(byte[] args){ - return NON_KEY_COMMANDS.contains(ByteBuffer.wrap(args)); - } - - private RedisOp parseFieldsCommand(RedisOpType opType, byte[][] args) { - int idx = keyStartIndex; - RedisKey key = new RedisKey(args[idx++]); - - // 跳过 FIELDS 之前的所有可选标志 - while (idx < args.length) { - ByteBuffer tokenBuf = ByteBuffer.wrap(args[idx]); - if (NON_KEY_COMMANDS.contains(tokenBuf)) { - idx++; - if (NON_KEY_COMMANDS_WITH_ARGS.contains(tokenBuf) && idx < args.length) { - idx++; // 跳过标志的参数 - } - } else if (Arrays.equals(args[idx], FIELDS_BYTES) || Arrays.equals(args[idx], FIELDS_BYTES_LOWER)) { - idx++; - break; - } else { - idx++; - } - } - - if (idx >= args.length) throw new IllegalArgumentException("Missing numfields"); - int numFields = Integer.parseInt(new String(args[idx++])); - List subKeys = new ArrayList<>(numFields); - - for (int i = 0; i < numFields; i++) { - if (idx >= args.length) throw new IllegalArgumentException("Incomplete field list"); - subKeys.add(new RedisKey(args[idx])); - idx += kvNum; - } - return new RedisOpMultiSubKey(opType, args, key, subKeys); + private boolean nonKey(byte[] arg) { + return NON_KEY_COMMANDS.contains(ByteBuffer.wrap(arg)); } -} +} \ No newline at end of file diff --git a/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/sequence/DefaultSequenceController.java b/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/sequence/DefaultSequenceController.java index 16f67ab787..31006ae957 100644 --- a/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/sequence/DefaultSequenceController.java +++ b/redis/redis-keeper/src/main/java/com/ctrip/xpipe/redis/keeper/applier/sequence/DefaultSequenceController.java @@ -325,8 +325,9 @@ private void submitMultiKeyCommand(RedisOpDataCommand command, long commandOf if(!multiRunningCommands.isEmpty()) { for (RedisKey key : keys) { - if (key.get()[0] == TAG_START) { - key = new RedisKey(client.hashTag(key.get())); + byte[] hashTag = client.hashTag(key.get()); + if (hashTag != null) { + key = new RedisKey(hashTag); } SequenceCommand lastSameKey = multiRunningCommands.get(key); if (lastSameKey != null) { @@ -368,8 +369,9 @@ private void submitObstacle(RedisOpCommand command, long commandOffset, GtidS RedisMultiKeyOp redisMultiKeyOp = (RedisMultiKeyOp) transactionOp; for(RedisKey redisKey:redisMultiKeyOp.getKeys()){ RedisKey commandKey = redisKey; - if (redisKey != null && redisKey.get()[0] == TAG_START) { - commandKey = new RedisKey(client.hashTag(redisKey.get())); + byte[] hashTag = client.hashTag(redisKey.get()); + if (hashTag != null) { + commandKey = new RedisKey(hashTag); transactionOpKeys.add(commandKey); break; }