diff --git a/example/subscription/src/main/java/org/apache/iotdb/ConsensusSubscriptionSessionExample.java b/example/subscription/src/main/java/org/apache/iotdb/ConsensusSubscriptionSessionExample.java index c0ebbe37198e8..12596f5d5b9d0 100644 --- a/example/subscription/src/main/java/org/apache/iotdb/ConsensusSubscriptionSessionExample.java +++ b/example/subscription/src/main/java/org/apache/iotdb/ConsensusSubscriptionSessionExample.java @@ -109,7 +109,7 @@ private static void createConsensusTopic(final String topicName, final String pa session.open(); final Properties config = new Properties(); - config.put(TopicConstant.MODE_KEY, TopicConstant.MODE_CONSENSUS_VALUE); + config.put(TopicConstant.MODE_KEY, TopicConstant.MODE_INCREMENTAL_VALUE); config.put(TopicConstant.FORMAT_KEY, TopicConstant.FORMAT_RECORD_HANDLER_VALUE); config.put(TopicConstant.PATH_KEY, path); config.put(TopicConstant.ORDER_MODE_KEY, TopicConstant.ORDER_MODE_PER_WRITER_VALUE); diff --git a/example/subscription/src/main/java/org/apache/iotdb/ConsensusTableModelSubscriptionSessionExample.java b/example/subscription/src/main/java/org/apache/iotdb/ConsensusTableModelSubscriptionSessionExample.java index a877a4a861eda..fa7733d57f4fe 100644 --- a/example/subscription/src/main/java/org/apache/iotdb/ConsensusTableModelSubscriptionSessionExample.java +++ b/example/subscription/src/main/java/org/apache/iotdb/ConsensusTableModelSubscriptionSessionExample.java @@ -107,7 +107,7 @@ private static void createConsensusTopic( .password(PASSWORD) .build()) { final Properties config = new Properties(); - config.put(TopicConstant.MODE_KEY, TopicConstant.MODE_CONSENSUS_VALUE); + config.put(TopicConstant.MODE_KEY, TopicConstant.MODE_INCREMENTAL_VALUE); config.put(TopicConstant.FORMAT_KEY, TopicConstant.FORMAT_RECORD_HANDLER_VALUE); config.put(TopicConstant.DATABASE_KEY, database); config.put(TopicConstant.TABLE_KEY, table); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/consensus/local/ConsensusSubscriptionITSupport.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/consensus/local/ConsensusSubscriptionITSupport.java index 254b5ffeb8558..066c38ef83145 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/consensus/local/ConsensusSubscriptionITSupport.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/consensus/local/ConsensusSubscriptionITSupport.java @@ -94,7 +94,7 @@ static void createConsensusTopic(final String topicName, final String path) thro session.dropTopicIfExists(topicName); final Properties config = new Properties(); - config.put(TopicConstant.MODE_KEY, TopicConstant.MODE_CONSENSUS_VALUE); + config.put(TopicConstant.MODE_KEY, TopicConstant.MODE_INCREMENTAL_VALUE); config.put(TopicConstant.FORMAT_KEY, TopicConstant.FORMAT_RECORD_HANDLER_VALUE); config.put(TopicConstant.PATH_KEY, path); config.put(TopicConstant.ORDER_MODE_KEY, TopicConstant.ORDER_MODE_PER_WRITER_VALUE); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/consensus/local/tablemodel/ConsensusSubscriptionTableITSupport.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/consensus/local/tablemodel/ConsensusSubscriptionTableITSupport.java index 4bd38992ed1f0..29a859d3011ee 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/consensus/local/tablemodel/ConsensusSubscriptionTableITSupport.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/consensus/local/tablemodel/ConsensusSubscriptionTableITSupport.java @@ -134,7 +134,7 @@ static void createConsensusTopic( session.dropTopicIfExists(topicName); final Properties config = new Properties(); - config.put(TopicConstant.MODE_KEY, TopicConstant.MODE_CONSENSUS_VALUE); + config.put(TopicConstant.MODE_KEY, TopicConstant.MODE_INCREMENTAL_VALUE); config.put(TopicConstant.FORMAT_KEY, TopicConstant.FORMAT_SESSION_DATA_SETS_HANDLER_VALUE); config.put(TopicConstant.DATABASE_KEY, databasePattern); config.put(TopicConstant.TABLE_KEY, tablePattern); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/consensus/local/tablemodel/IoTDBConsensusSubscriptionColumnFilterClusterIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/consensus/local/tablemodel/IoTDBConsensusSubscriptionColumnFilterClusterIT.java index d2fd5552ee75e..6c33532828b09 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/consensus/local/tablemodel/IoTDBConsensusSubscriptionColumnFilterClusterIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/consensus/local/tablemodel/IoTDBConsensusSubscriptionColumnFilterClusterIT.java @@ -150,7 +150,7 @@ private static void createOwnedConsensusTopic( session.dropTopicIfExists(topicName); final Properties config = new Properties(); - config.put(TopicConstant.MODE_KEY, TopicConstant.MODE_CONSENSUS_VALUE); + config.put(TopicConstant.MODE_KEY, TopicConstant.MODE_INCREMENTAL_VALUE); config.put(TopicConstant.FORMAT_KEY, TopicConstant.FORMAT_SESSION_DATA_SETS_HANDLER_VALUE); config.put(TopicConstant.DATABASE_KEY, database); config.put(TopicConstant.TABLE_KEY, table); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/dual/tablemodel/IoTDBSubscriptionColumnFilterIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/dual/tablemodel/IoTDBSubscriptionColumnFilterIT.java index 27cd10ed5340b..bfeb7e6b7aa7b 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/dual/tablemodel/IoTDBSubscriptionColumnFilterIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/dual/tablemodel/IoTDBSubscriptionColumnFilterIT.java @@ -1015,7 +1015,7 @@ private void createTopic( final String columnFilter) throws Exception { createTopic( - topicName, database, tableName, TopicConstant.MODE_LIVE_VALUE, format, columnFilter); + topicName, database, tableName, TopicConstant.MODE_INITIAL_VALUE, format, columnFilter); } private void createTopic( @@ -1032,7 +1032,8 @@ private void createTopic( private void createTopicWithoutColumnFilter( final String topicName, final String database, final String tableName, final String format) throws Exception { - createTopic(topicName, database, tableName, TopicConstant.MODE_LIVE_VALUE, format, "", false); + createTopic( + topicName, database, tableName, TopicConstant.MODE_INITIAL_VALUE, format, "", false); } private void createTopic( diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/local/tablemodel/IoTDBSubscriptionPermissionIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/local/tablemodel/IoTDBSubscriptionPermissionIT.java index 123a68c8df53e..44d0dea16d77f 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/local/tablemodel/IoTDBSubscriptionPermissionIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/local/tablemodel/IoTDBSubscriptionPermissionIT.java @@ -390,7 +390,7 @@ public void testColumnFilterTopicRequiresSystemPermission() throws Exception { private static Properties columnFilterTopicConfig( final String database, final String tableName, final String columnFilter) { final Properties topicConfig = new Properties(); - topicConfig.put(TopicConstant.MODE_KEY, TopicConstant.MODE_LIVE_VALUE); + topicConfig.put(TopicConstant.MODE_KEY, TopicConstant.MODE_INITIAL_VALUE); topicConfig.put(TopicConstant.FORMAT_KEY, TopicConstant.FORMAT_RECORD_HANDLER_VALUE); topicConfig.put(TopicConstant.DATABASE_KEY, database); topicConfig.put(TopicConstant.TABLE_KEY, tableName); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBAllTsDatasetPullConsumerIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBAllTsDatasetPullConsumerIT.java index a1a0cae5d9cce..fe2ac71aba773 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBAllTsDatasetPullConsumerIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBAllTsDatasetPullConsumerIT.java @@ -75,7 +75,7 @@ public void setUp() throws Exception { "2024-01-01T00:00:00+08:00", "2024-03-31T23:59:59+08:00", false, - TopicConstant.MODE_LIVE_VALUE, + TopicConstant.MODE_INITIAL_VALUE, TopicConstant.LOOSE_RANGE_ALL_VALUE); session_src.createTimeseries( pattern, TSDataType.INT64, TSEncoding.GORILLA, CompressionType.LZ4); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBAllTsTsfilePullConsumerIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBAllTsTsfilePullConsumerIT.java index c00ea2401ce40..8a585c5727e2c 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBAllTsTsfilePullConsumerIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBAllTsTsfilePullConsumerIT.java @@ -78,7 +78,7 @@ public void setUp() throws Exception { null, String.valueOf(nowTimestamp), true, - TopicConstant.MODE_LIVE_VALUE, + TopicConstant.MODE_INITIAL_VALUE, TopicConstant.LOOSE_RANGE_ALL_VALUE); session_src.createTimeseries( device + ".s_0", TSDataType.INT64, TSEncoding.GORILLA, CompressionType.LZ4); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBPathDeviceDataSetPullConsumerIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBPathDeviceDataSetPullConsumerIT.java index 2cb4b757253ef..a048041ce90da 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBPathDeviceDataSetPullConsumerIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBPathDeviceDataSetPullConsumerIT.java @@ -74,7 +74,7 @@ public void setUp() throws Exception { null, "now", false, - TopicConstant.MODE_LIVE_VALUE, + TopicConstant.MODE_INITIAL_VALUE, TopicConstant.LOOSE_RANGE_PATH_VALUE); session_src.createTimeseries( device + ".s_0", TSDataType.INT64, TSEncoding.GORILLA, CompressionType.LZ4); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBPathDeviceTsfilePullConsumerIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBPathDeviceTsfilePullConsumerIT.java index 70a5857ac8851..746e2abdcb461 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBPathDeviceTsfilePullConsumerIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBPathDeviceTsfilePullConsumerIT.java @@ -76,7 +76,7 @@ public void setUp() throws Exception { null, null, true, - TopicConstant.MODE_LIVE_VALUE, + TopicConstant.MODE_INITIAL_VALUE, TopicConstant.LOOSE_RANGE_PATH_VALUE); session_src.createTimeseries( device + ".s_0", TSDataType.INT64, TSEncoding.GORILLA, CompressionType.LZ4); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBTimeTsDatasetPullConsumerIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBTimeTsDatasetPullConsumerIT.java index 3fc35e39bb9b4..1b8e9174e9bf3 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBTimeTsDatasetPullConsumerIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBTimeTsDatasetPullConsumerIT.java @@ -76,7 +76,7 @@ public void setUp() throws Exception { "2024-01-01T00:00:00+08:00", "2024-03-31T23:59:59+08:00", false, - TopicConstant.MODE_LIVE_VALUE, + TopicConstant.MODE_INITIAL_VALUE, TopicConstant.LOOSE_RANGE_TIME_VALUE); session_src.createTimeseries( pattern, TSDataType.INT64, TSEncoding.GORILLA, CompressionType.LZ4); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBTimeTsTsfilePullConsumerIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBTimeTsTsfilePullConsumerIT.java index 7cec98a627015..f609702e43145 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBTimeTsTsfilePullConsumerIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pullconsumer/loose_range/IoTDBTimeTsTsfilePullConsumerIT.java @@ -79,7 +79,7 @@ public void setUp() throws Exception { null, String.valueOf(nowTimestamp), true, - TopicConstant.MODE_LIVE_VALUE, + TopicConstant.MODE_INITIAL_VALUE, TopicConstant.LOOSE_RANGE_TIME_VALUE); session_src.createTimeseries( device + ".s_0", TSDataType.INT64, TSEncoding.GORILLA, CompressionType.LZ4); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBLooseAllTsDatasetPushConsumerIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBLooseAllTsDatasetPushConsumerIT.java index c30e9d8493dab..af8d620637fdc 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBLooseAllTsDatasetPushConsumerIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBLooseAllTsDatasetPushConsumerIT.java @@ -56,7 +56,7 @@ * DataSet * pattern: ts * loose-range: all - * mode: live + * mode: initial */ @RunWith(IoTDBTestRunner.class) @Category({MultiClusterIT2SubscriptionTreeRegressionConsumer.class}) @@ -82,7 +82,7 @@ public void setUp() throws Exception { "2024-01-01T00:00:00+08:00", "2024-02-13T08:00:02+08:00", false, - TopicConstant.MODE_LIVE_VALUE, + TopicConstant.MODE_INITIAL_VALUE, TopicConstant.LOOSE_RANGE_ALL_VALUE); session_src.createTimeseries( device + ".s_0", TSDataType.INT64, TSEncoding.GORILLA, CompressionType.LZ4); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBLooseAllTsfilePushConsumerIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBLooseAllTsfilePushConsumerIT.java index 91fb1a281e501..1c2f420f339e6 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBLooseAllTsfilePushConsumerIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBLooseAllTsfilePushConsumerIT.java @@ -58,7 +58,7 @@ /*** * push consumer - * mode: live + * mode: initial * pattern: db * loose-range: all */ @@ -84,7 +84,7 @@ public void setUp() throws Exception { "2024-01-01T00:00:00+08:00", "2024-03-31T00:00:00+08:00", true, - TopicConstant.MODE_LIVE_VALUE, + TopicConstant.MODE_INITIAL_VALUE, TopicConstant.LOOSE_RANGE_ALL_VALUE); session_src.createTimeseries( device + ".s_0", TSDataType.INT64, TSEncoding.GORILLA, CompressionType.LZ4); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBPathLooseDeviceTsfilePushConsumerIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBPathLooseDeviceTsfilePushConsumerIT.java index c647160493a29..270d5cd928f94 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBPathLooseDeviceTsfilePushConsumerIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBPathLooseDeviceTsfilePushConsumerIT.java @@ -83,7 +83,7 @@ public void setUp() throws Exception { "2024-01-01T00:00:00+08:00", "2024-03-31T00:00:00+08:00", true, - TopicConstant.MODE_LIVE_VALUE, + TopicConstant.MODE_INITIAL_VALUE, TopicConstant.LOOSE_RANGE_PATH_VALUE); session_src.createTimeseries( device + ".s_0", TSDataType.INT64, TSEncoding.GORILLA, CompressionType.LZ4); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBPathLooseTsDatasetPushConsumerIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBPathLooseTsDatasetPushConsumerIT.java index 3e144c137fd3b..de854f29317dd 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBPathLooseTsDatasetPushConsumerIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBPathLooseTsDatasetPushConsumerIT.java @@ -83,7 +83,7 @@ public void setUp() throws Exception { "2024-01-01T00:00:00+08:00", "2024-02-13T08:00:02+08:00", false, - TopicConstant.MODE_LIVE_VALUE, + TopicConstant.MODE_INITIAL_VALUE, TopicConstant.LOOSE_RANGE_PATH_VALUE); session_src.createTimeseries( device + ".s_0", TSDataType.INT64, TSEncoding.GORILLA, CompressionType.LZ4); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBPathLooseTsfilePushConsumerIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBPathLooseTsfilePushConsumerIT.java index b8c438bfb9d19..c8800aa93ac25 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBPathLooseTsfilePushConsumerIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBPathLooseTsfilePushConsumerIT.java @@ -80,7 +80,7 @@ public void setUp() throws Exception { "2024-01-01T00:00:00+08:00", "2024-03-31T00:00:00+08:00", true, - TopicConstant.MODE_LIVE_VALUE, + TopicConstant.MODE_INITIAL_VALUE, TopicConstant.LOOSE_RANGE_PATH_VALUE); session_src.createTimeseries( device + ".s_0", TSDataType.INT64, TSEncoding.GORILLA, CompressionType.LZ4); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBPathTsLooseDatasetPushConsumerIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBPathTsLooseDatasetPushConsumerIT.java index 5a2a5d6a35110..aba24718a632d 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBPathTsLooseDatasetPushConsumerIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBPathTsLooseDatasetPushConsumerIT.java @@ -81,7 +81,7 @@ public void setUp() throws Exception { "2024-01-01T00:00:00+08:00", "2024-02-13T08:00:02+08:00", false, - TopicConstant.MODE_LIVE_VALUE, + TopicConstant.MODE_INITIAL_VALUE, TopicConstant.LOOSE_RANGE_PATH_VALUE); session_src.createTimeseries( device + ".s_0", TSDataType.INT64, TSEncoding.GORILLA, CompressionType.LZ4); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBTimeLooseTsDatasetPushConsumerIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBTimeLooseTsDatasetPushConsumerIT.java index ecbd661e5e076..38b826f3cafc5 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBTimeLooseTsDatasetPushConsumerIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBTimeLooseTsDatasetPushConsumerIT.java @@ -57,7 +57,7 @@ * DataSet * pattern: ts * loose-range: time - * live + * initial */ @RunWith(IoTDBTestRunner.class) @Category({MultiClusterIT2SubscriptionTreeRegressionConsumer.class}) @@ -84,7 +84,7 @@ public void setUp() throws Exception { "2024-01-01T00:00:00+08:00", "2024-02-13T08:00:02+08:00", false, - TopicConstant.MODE_LIVE_VALUE, + TopicConstant.MODE_INITIAL_VALUE, TopicConstant.LOOSE_RANGE_TIME_VALUE); session_src.createTimeseries( device + ".s_0", TSDataType.INT64, TSEncoding.GORILLA, CompressionType.LZ4); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBTimeLooseTsTsfilePushConsumerIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBTimeLooseTsTsfilePushConsumerIT.java index 0311bc0a0b561..e62644a59dca3 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBTimeLooseTsTsfilePushConsumerIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBTimeLooseTsTsfilePushConsumerIT.java @@ -57,7 +57,7 @@ import static org.apache.iotdb.subscription.it.IoTDBSubscriptionITConstant.AWAIT; /*** - * mode: live + * mode: initial * loose-range:path * format: tsfile */ @@ -83,7 +83,7 @@ public void setUp() throws Exception { "2024-01-01T00:00:00+08:00", "2024-03-31T00:00:00+08:00", true, - TopicConstant.MODE_LIVE_VALUE, + TopicConstant.MODE_INITIAL_VALUE, TopicConstant.LOOSE_RANGE_TIME_VALUE); session_src.createTimeseries( device + ".s_0", TSDataType.INT64, TSEncoding.GORILLA, CompressionType.LZ4); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBTimeLooseTsfilePushConsumerIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBTimeLooseTsfilePushConsumerIT.java index 53207ee7fe633..68b36104b016f 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBTimeLooseTsfilePushConsumerIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBTimeLooseTsfilePushConsumerIT.java @@ -77,7 +77,7 @@ public void setUp() throws Exception { "2024-01-01T00:00:00+08:00", "2024-03-31T00:00:00+08:00", true, - TopicConstant.MODE_LIVE_VALUE, + TopicConstant.MODE_INITIAL_VALUE, TopicConstant.LOOSE_RANGE_TIME_VALUE); session_src.createTimeseries( device + ".s_0", TSDataType.INT64, TSEncoding.GORILLA, CompressionType.LZ4); diff --git a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBTimeTsLooseDatasetPushConsumerIT.java b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBTimeTsLooseDatasetPushConsumerIT.java index 4f358ddb8c464..bfd539054e9b3 100644 --- a/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBTimeTsLooseDatasetPushConsumerIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/subscription/it/triple/treemodel/regression/pushconsumer/loose_range/IoTDBTimeTsLooseDatasetPushConsumerIT.java @@ -56,7 +56,7 @@ * DataSet * pattern: ts * time loose - * live + * initial */ @RunWith(IoTDBTestRunner.class) @Category({MultiClusterIT2SubscriptionTreeRegressionConsumer.class}) @@ -83,7 +83,7 @@ public void setUp() throws Exception { "2024-01-01T00:00:00+08:00", "2024-02-13T08:00:02+08:00", false, - TopicConstant.MODE_LIVE_VALUE, + TopicConstant.MODE_INITIAL_VALUE, TopicConstant.LOOSE_RANGE_TIME_VALUE); session_src.createTimeseries( device + ".s_0", TSDataType.INT64, TSEncoding.GORILLA, CompressionType.LZ4); diff --git a/iotdb-client/subscription/src/main/i18n/en/org/apache/iotdb/rpc/subscription/i18n/SubscriptionMessages.java b/iotdb-client/subscription/src/main/i18n/en/org/apache/iotdb/rpc/subscription/i18n/SubscriptionMessages.java index 7ce7e7a876381..95f26d55d681c 100644 --- a/iotdb-client/subscription/src/main/i18n/en/org/apache/iotdb/rpc/subscription/i18n/SubscriptionMessages.java +++ b/iotdb-client/subscription/src/main/i18n/en/org/apache/iotdb/rpc/subscription/i18n/SubscriptionMessages.java @@ -270,7 +270,7 @@ private SubscriptionMessages() {} public static final String EXCEPTION_CLUSTER_HAS_NO_AVAILABLE_SUBSCRIPTION_PROVIDERS_ARG_FETCH_ALL_ENDPOINTS_D232693E = "Cluster has no available subscription providers when %s fetch all endpoints"; public static final String EXCEPTION_PASSWORD_ENCRYPTEDPASSWORD_MUTUALLY_EXCLUSIVE_ENCRYPTEDPASSWORD_ALREADY_SET_E4548A43 = "password and encryptedPassword are mutually exclusive; encryptedPassword is already set"; public static final String EXCEPTION_PASSWORD_ENCRYPTEDPASSWORD_MUTUALLY_EXCLUSIVE_PASSWORD_ALREADY_SET_BB20AD1E = "password and encryptedPassword are mutually exclusive; password is already set"; - public static final String EXCEPTION_CONSENSUS_MODE_TOPIC_SHOULD_NOT_GENERATE_PIPE_SOURCE_ATTRIBUTES_BBDFF732 = "Consensus mode topic should not generate pipe source attributes"; + public static final String EXCEPTION_INCREMENTAL_MODE_TOPIC_SHOULD_NOT_GENERATE_PIPE_SOURCE_ATTRIBUTES_09D75393 = "Incremental mode topic should not generate pipe source attributes"; public static final String EXCEPTION_UNSUPPORTED_SUBSCRIPTIONCOMMITCONTEXT_VERSION_8021B27B = "Unsupported SubscriptionCommitContext version: "; public static final String OUTDATED_SUBSCRIPTION_EVENT = "outdated subscription event"; public static final String FIELD_SEPARATOR = ", "; diff --git a/iotdb-client/subscription/src/main/i18n/zh/org/apache/iotdb/rpc/subscription/i18n/SubscriptionMessages.java b/iotdb-client/subscription/src/main/i18n/zh/org/apache/iotdb/rpc/subscription/i18n/SubscriptionMessages.java index 4a8edd3f54a85..562a4d7257281 100644 --- a/iotdb-client/subscription/src/main/i18n/zh/org/apache/iotdb/rpc/subscription/i18n/SubscriptionMessages.java +++ b/iotdb-client/subscription/src/main/i18n/zh/org/apache/iotdb/rpc/subscription/i18n/SubscriptionMessages.java @@ -219,7 +219,7 @@ private SubscriptionMessages() {} public static final String EXCEPTION_CLUSTER_HAS_NO_AVAILABLE_SUBSCRIPTION_PROVIDERS_ARG_FETCH_ALL_ENDPOINTS_D232693E = "%s 获取所有 endpoint 时,集群没有可用的 SubscriptionProvider"; public static final String EXCEPTION_PASSWORD_ENCRYPTEDPASSWORD_MUTUALLY_EXCLUSIVE_ENCRYPTEDPASSWORD_ALREADY_SET_E4548A43 = "password 与 encryptedPassword 互斥;已设置 encryptedPassword"; public static final String EXCEPTION_PASSWORD_ENCRYPTEDPASSWORD_MUTUALLY_EXCLUSIVE_PASSWORD_ALREADY_SET_BB20AD1E = "password 与 encryptedPassword 互斥;已设置 password"; - public static final String EXCEPTION_CONSENSUS_MODE_TOPIC_SHOULD_NOT_GENERATE_PIPE_SOURCE_ATTRIBUTES_BBDFF732 = "Consensus mode 主题不应生成 pipe source attributes"; + public static final String EXCEPTION_INCREMENTAL_MODE_TOPIC_SHOULD_NOT_GENERATE_PIPE_SOURCE_ATTRIBUTES_09D75393 = "incremental mode 的 topic 不应生成 pipe source attributes"; public static final String EXCEPTION_UNSUPPORTED_SUBSCRIPTIONCOMMITCONTEXT_VERSION_8021B27B = "不支持的 SubscriptionCommitContext 版本:"; public static final String OUTDATED_SUBSCRIPTION_EVENT = "过期的订阅事件"; public static final String FIELD_SEPARATOR = ","; diff --git a/iotdb-client/subscription/src/main/java/org/apache/iotdb/rpc/subscription/config/TopicConfig.java b/iotdb-client/subscription/src/main/java/org/apache/iotdb/rpc/subscription/config/TopicConfig.java index b0824fcb39427..3d107476b85b1 100644 --- a/iotdb-client/subscription/src/main/java/org/apache/iotdb/rpc/subscription/config/TopicConfig.java +++ b/iotdb-client/subscription/src/main/java/org/apache/iotdb/rpc/subscription/config/TopicConfig.java @@ -30,6 +30,7 @@ import java.util.Collections; import java.util.HashMap; import java.util.HashSet; +import java.util.Locale; import java.util.Map; import java.util.Objects; import java.util.Set; @@ -43,8 +44,8 @@ public class TopicConfig extends PipeParameters { static { final Set modes = new HashSet<>(3); modes.add(TopicConstant.MODE_SNAPSHOT_VALUE); - modes.add(TopicConstant.MODE_LIVE_VALUE); - modes.add(TopicConstant.MODE_CONSENSUS_VALUE); + modes.add(TopicConstant.MODE_INITIAL_VALUE); + modes.add(TopicConstant.MODE_INCREMENTAL_VALUE); MODE_VALUE_SET = Collections.unmodifiableSet(modes); final Set orderModes = new HashSet<>(3); @@ -84,8 +85,10 @@ public TopicConfig(final Map attributes) { private static final Map SNAPSHOT_MODE_CONFIG = Collections.singletonMap("mode", TopicConstant.MODE_SNAPSHOT_VALUE); - private static final Map LIVE_MODE_CONFIG = - Collections.singletonMap("mode", TopicConstant.MODE_LIVE_VALUE); + // Pipe source still uses the legacy "live" value for the initial (full + incremental) topic + // mode. + private static final Map INITIAL_MODE_CONFIG = + Collections.singletonMap("mode", TopicConstant.LEGACY_MODE_LIVE_VALUE); private static final Map STRICT_MODE_CONFIG = Collections.singletonMap("mode.strict", "true"); @@ -125,12 +128,28 @@ public boolean isSnapshotMode() { return TopicConstant.MODE_SNAPSHOT_VALUE.equalsIgnoreCase(getMode()); } + public boolean isInitialMode() { + return TopicConstant.MODE_INITIAL_VALUE.equalsIgnoreCase(getMode()); + } + + public boolean isIncrementalMode() { + return TopicConstant.MODE_INCREMENTAL_VALUE.equalsIgnoreCase(getMode()); + } + + /** + * @deprecated Use {@link #isInitialMode()}. + */ + @Deprecated public boolean isLiveMode() { - return TopicConstant.MODE_LIVE_VALUE.equalsIgnoreCase(getMode()); + return isInitialMode(); } + /** + * @deprecated Use {@link #isIncrementalMode()}. + */ + @Deprecated public boolean isConsensusMode() { - return TopicConstant.MODE_CONSENSUS_VALUE.equalsIgnoreCase(getMode()); + return isIncrementalMode(); } public static boolean isValidMode(final String mode) { @@ -138,7 +157,19 @@ public static boolean isValidMode(final String mode) { } public static String normalizeMode(final String mode) { - return mode == null ? TopicConstant.MODE_DEFAULT_VALUE : mode.trim().toLowerCase(); + if (mode == null) { + return TopicConstant.MODE_DEFAULT_VALUE; + } + + final String normalizedMode = mode.trim().toLowerCase(Locale.ROOT); + switch (normalizedMode) { + case TopicConstant.LEGACY_MODE_LIVE_VALUE: + return TopicConstant.MODE_INITIAL_VALUE; + case TopicConstant.LEGACY_MODE_CONSENSUS_VALUE: + return TopicConstant.MODE_INCREMENTAL_VALUE; + default: + return normalizedMode; + } } public String getOrderMode() { @@ -223,12 +254,12 @@ public Map getAttributesWithSourceRealtimeMode() { } public Map getAttributesWithSourceMode() { - if (isConsensusMode()) { + if (isIncrementalMode()) { throw new IllegalArgumentException( SubscriptionMessages - .EXCEPTION_CONSENSUS_MODE_TOPIC_SHOULD_NOT_GENERATE_PIPE_SOURCE_ATTRIBUTES_BBDFF732); + .EXCEPTION_INCREMENTAL_MODE_TOPIC_SHOULD_NOT_GENERATE_PIPE_SOURCE_ATTRIBUTES_09D75393); } - return isSnapshotMode() ? SNAPSHOT_MODE_CONFIG : LIVE_MODE_CONFIG; + return isSnapshotMode() ? SNAPSHOT_MODE_CONFIG : INITIAL_MODE_CONFIG; } public Map getAttributesWithSourceLooseRangeOrStrict() { diff --git a/iotdb-client/subscription/src/main/java/org/apache/iotdb/rpc/subscription/config/TopicConstant.java b/iotdb-client/subscription/src/main/java/org/apache/iotdb/rpc/subscription/config/TopicConstant.java index f929d7b74722e..4d11f2fe8a27a 100644 --- a/iotdb-client/subscription/src/main/java/org/apache/iotdb/rpc/subscription/config/TopicConstant.java +++ b/iotdb-client/subscription/src/main/java/org/apache/iotdb/rpc/subscription/config/TopicConstant.java @@ -42,10 +42,23 @@ public class TopicConstant { public static final String NOW_TIME_VALUE = "now"; public static final String MODE_KEY = "mode"; - public static final String MODE_LIVE_VALUE = "live"; + public static final String MODE_INITIAL_VALUE = "initial"; public static final String MODE_SNAPSHOT_VALUE = "snapshot"; - public static final String MODE_CONSENSUS_VALUE = "consensus"; - public static final String MODE_DEFAULT_VALUE = MODE_LIVE_VALUE; + public static final String MODE_INCREMENTAL_VALUE = "incremental"; + public static final String MODE_DEFAULT_VALUE = MODE_INITIAL_VALUE; + + static final String LEGACY_MODE_LIVE_VALUE = "live"; + static final String LEGACY_MODE_CONSENSUS_VALUE = "consensus"; + + /** + * @deprecated Use {@link #MODE_INITIAL_VALUE}. + */ + @Deprecated public static final String MODE_LIVE_VALUE = LEGACY_MODE_LIVE_VALUE; + + /** + * @deprecated Use {@link #MODE_INCREMENTAL_VALUE}. + */ + @Deprecated public static final String MODE_CONSENSUS_VALUE = LEGACY_MODE_CONSENSUS_VALUE; public static final String ORDER_MODE_KEY = "order-mode"; public static final String ORDER_MODE_LEADER_ONLY_VALUE = "leader-only"; diff --git a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumer.java b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumer.java index 74491898e5948..fa48b6a099a9f 100644 --- a/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumer.java +++ b/iotdb-client/subscription/src/main/java/org/apache/iotdb/session/subscription/consumer/base/AbstractSubscriptionConsumer.java @@ -168,7 +168,7 @@ public boolean allTopicMessagesHaveBeenConsumed() { private boolean allTopicMessagesHaveBeenConsumed(final Collection topicNames) { // For the topic that needs to be detected, there are two scenarios to consider: - // 1. If configs as live, it cannot be determined whether the topic has been fully consumed. + // 1. Initial topics are unbounded and cannot be fully consumed. // 2. If configs as snapshot, it means the topic has not been automatically unsubscribed. // Therefore, the logic can be summarized as follows: if there is a matching topic in subscribed // topics, then it has not been fully consumed. diff --git a/iotdb-client/subscription/src/test/java/org/apache/iotdb/rpc/subscription/config/TopicConfigTest.java b/iotdb-client/subscription/src/test/java/org/apache/iotdb/rpc/subscription/config/TopicConfigTest.java index d60f4f1f7021a..d308ef123cf68 100644 --- a/iotdb-client/subscription/src/test/java/org/apache/iotdb/rpc/subscription/config/TopicConfigTest.java +++ b/iotdb-client/subscription/src/test/java/org/apache/iotdb/rpc/subscription/config/TopicConfigTest.java @@ -28,6 +28,59 @@ public class TopicConfigTest { + @Test + public void testModeDefaultsToInitial() { + final TopicConfig topicConfig = new TopicConfig(); + + Assert.assertEquals(TopicConstant.MODE_INITIAL_VALUE, topicConfig.getMode()); + Assert.assertTrue(topicConfig.isInitialMode()); + Assert.assertFalse(topicConfig.isSnapshotMode()); + Assert.assertFalse(topicConfig.isIncrementalMode()); + } + + @Test + public void testCanonicalModeValues() { + Assert.assertTrue(TopicConfig.isValidMode(TopicConstant.MODE_INITIAL_VALUE)); + Assert.assertTrue(TopicConfig.isValidMode(TopicConstant.MODE_SNAPSHOT_VALUE)); + Assert.assertTrue(TopicConfig.isValidMode(TopicConstant.MODE_INCREMENTAL_VALUE)); + Assert.assertFalse(TopicConfig.isValidMode("wal")); + + Assert.assertTrue(topicConfigWithMode(" INITIAL ").isInitialMode()); + Assert.assertTrue(topicConfigWithMode(" INCREMENTAL ").isIncrementalMode()); + } + + @SuppressWarnings("deprecation") + @Test + public void testLegacyModeValues() { + final TopicConfig liveTopicConfig = topicConfigWithMode(TopicConstant.MODE_LIVE_VALUE); + Assert.assertTrue(TopicConfig.isValidMode(TopicConstant.MODE_LIVE_VALUE)); + Assert.assertEquals(TopicConstant.MODE_INITIAL_VALUE, liveTopicConfig.getMode()); + Assert.assertTrue(liveTopicConfig.isInitialMode()); + Assert.assertTrue(liveTopicConfig.isLiveMode()); + + final TopicConfig consensusTopicConfig = + topicConfigWithMode(TopicConstant.MODE_CONSENSUS_VALUE); + Assert.assertTrue(TopicConfig.isValidMode(TopicConstant.MODE_CONSENSUS_VALUE)); + Assert.assertEquals(TopicConstant.MODE_INCREMENTAL_VALUE, consensusTopicConfig.getMode()); + Assert.assertTrue(consensusTopicConfig.isIncrementalMode()); + Assert.assertTrue(consensusTopicConfig.isConsensusMode()); + } + + @SuppressWarnings("deprecation") + @Test + public void testInitialModeMapsToPipeLiveMode() { + Assert.assertEquals( + TopicConstant.MODE_LIVE_VALUE, + topicConfigWithMode(TopicConstant.MODE_INITIAL_VALUE) + .getAttributesWithSourceMode() + .get(TopicConstant.MODE_KEY)); + Assert.assertEquals( + TopicConstant.MODE_LIVE_VALUE, + topicConfigWithMode(TopicConstant.MODE_LIVE_VALUE) + .getAttributesWithSourceMode() + .get(TopicConstant.MODE_KEY)); + } + @Test public void testColumnFilterKeyIsCaseInsensitive() { final TopicConfig topicConfig = @@ -56,4 +109,8 @@ public void testColumnFilterTrivialWithMixedCaseKeyAndValue() { Assert.assertTrue(new TopicConfig(attributes).isColumnFilterTrivial()); } + + private static TopicConfig topicConfigWithMode(final String mode) { + return new TopicConfig(Collections.singletonMap(TopicConstant.MODE_KEY, mode)); + } } diff --git a/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ConfigNodeMessages.java b/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ConfigNodeMessages.java index 014ad8e807c8d..54e06e9a97228 100644 --- a/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ConfigNodeMessages.java +++ b/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ConfigNodeMessages.java @@ -671,6 +671,9 @@ private ConfigNodeMessages() {} "procedure_completed_evict_ttl should be greater than 0, but was "; public static final String - EXCEPTION_FAILED_TO_CREATE_OR_ALTER_TOPIC_MODE_CONSENSUS_DOES_NOT_SUPPORT_TOPIC_ATTRIBUTES_ARG_3C2D0BDA = - "Failed to create or alter topic, mode=consensus does not support topic attributes %s"; + EXCEPTION_FAILED_TO_CREATE_OR_ALTER_TOPIC_MODE_INCREMENTAL_DOES_NOT_SUPPORT_TOPIC_ATTRIBUTES_ARG_1A72326A = + "Failed to create or alter topic, mode=incremental does not support topic attributes %s"; + public static final String + EXCEPTION_FAILED_TO_CREATE_OR_ALTER_TOPIC_ARG_AND_ARG_ARE_ONLY_SUPPORTED_FOR_INCREMENTAL_TOPICS_D86CEA8E = + "Failed to create or alter topic, %s and %s are only supported for incremental topics"; } diff --git a/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ConfigNodeMessages.java b/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ConfigNodeMessages.java index 3c3bce0b84575..0c28efc409677 100644 --- a/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ConfigNodeMessages.java +++ b/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ConfigNodeMessages.java @@ -716,6 +716,9 @@ private ConfigNodeMessages() {} EXCEPTION_PROCEDURE_COMPLETED_EVICT_TTL_SHOULD_BE_GREATER_THAN_0_BUT_WAS_5A4D0CF6 = "procedure_completed_evict_ttl 应大于 0,但当前值为 "; public static final String - EXCEPTION_FAILED_TO_CREATE_OR_ALTER_TOPIC_MODE_CONSENSUS_DOES_NOT_SUPPORT_TOPIC_ATTRIBUTES_ARG_3C2D0BDA = - "创建或修改 topic 失败,mode=consensus 不支持 topic 属性 %s"; + EXCEPTION_FAILED_TO_CREATE_OR_ALTER_TOPIC_MODE_INCREMENTAL_DOES_NOT_SUPPORT_TOPIC_ATTRIBUTES_ARG_1A72326A = + "创建或修改 topic 失败,mode=incremental 不支持 topic 属性 %s"; + public static final String + EXCEPTION_FAILED_TO_CREATE_OR_ALTER_TOPIC_ARG_AND_ARG_ARE_ONLY_SUPPORTED_FOR_INCREMENTAL_TOPICS_D86CEA8E = + "创建或修改 topic 失败,%s 和 %s 仅支持 incremental 模式的 topic"; } diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/subscription/runtime/SubscriptionRuntimeCoordinator.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/subscription/runtime/SubscriptionRuntimeCoordinator.java index 1f3a70073c5a5..dd8fe2e70cf7e 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/subscription/runtime/SubscriptionRuntimeCoordinator.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/subscription/runtime/SubscriptionRuntimeCoordinator.java @@ -102,7 +102,7 @@ public boolean hasAnyConsensusBasedTopic() { .getSubscriptionCoordinator() .getSubscriptionInfo() .getAllTopicMeta()) { - if (topicMeta.getConfig().isConsensusMode()) { + if (topicMeta.getConfig().isIncrementalMode()) { return true; } } diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/subscription/SubscriptionInfo.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/subscription/SubscriptionInfo.java index 4d063e83152b9..12b25159fa9bc 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/subscription/SubscriptionInfo.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/subscription/SubscriptionInfo.java @@ -110,7 +110,7 @@ public class SubscriptionInfo implements SnapshotProcessor { TopicConstant.OWNER_EPOCH_KEY, TopicConstant.MAX_OWNER_EPOCH_KEY, TopicConstant.OWNER_LEASE_DURATION_MS_KEY); - private static final Set CONSENSUS_TOPIC_SUPPORTED_ATTRIBUTE_KEYS = + private static final Set INCREMENTAL_TOPIC_SUPPORTED_ATTRIBUTE_KEYS = Set.of( SystemConstant.SQL_DIALECT_KEY, TopicConstant.PATH_KEY, @@ -340,21 +340,21 @@ private void validateTopicConfig(final TopicConfig topicConfig) throws Subscript TopicConstant.MODE_KEY, mode, TopicConstant.MODE_SNAPSHOT_VALUE, - TopicConstant.MODE_LIVE_VALUE, - TopicConstant.MODE_CONSENSUS_VALUE); + TopicConstant.MODE_INITIAL_VALUE, + TopicConstant.MODE_INCREMENTAL_VALUE); LOGGER.warn(exceptionMessage); throw new SubscriptionException(exceptionMessage); } - validateConsensusTopicAttributes(topicConfig); - validateConsensusProtocolSupport(topicConfig); + validateIncrementalTopicAttributes(topicConfig); + validateIncrementalProtocolSupport(topicConfig); - if (topicConfig.isConsensusMode() && !topicConfig.isRecordFormat()) { + if (topicConfig.isIncrementalMode() && !topicConfig.isRecordFormat()) { final String exceptionMessage = String.format( "Failed to create or alter topic, %s=%s only supports %s=%s", TopicConstant.MODE_KEY, - TopicConstant.MODE_CONSENSUS_VALUE, + TopicConstant.MODE_INCREMENTAL_VALUE, TopicConstant.FORMAT_KEY, TopicConstant.FORMAT_RECORD_HANDLER_VALUE); LOGGER.warn(exceptionMessage); @@ -376,7 +376,7 @@ private void validateTopicConfig(final TopicConfig topicConfig) throws Subscript } validateColumnFilter(topicConfig); - validateConsensusTopicRetentionConfig(topicConfig); + validateIncrementalTopicRetentionConfig(topicConfig); final Long ownerLeaseDurationMs = topicConfig.getLong(TopicConstant.OWNER_LEASE_DURATION_MS_KEY); @@ -393,9 +393,9 @@ private void validateTopicConfig(final TopicConfig topicConfig) throws Subscript } } - private void validateConsensusTopicAttributes(final TopicConfig topicConfig) + private void validateIncrementalTopicAttributes(final TopicConfig topicConfig) throws SubscriptionException { - if (!topicConfig.isConsensusMode()) { + if (!topicConfig.isIncrementalMode()) { return; } @@ -404,7 +404,7 @@ private void validateConsensusTopicAttributes(final TopicConfig topicConfig) .filter( key -> Objects.isNull(key) - || !CONSENSUS_TOPIC_SUPPORTED_ATTRIBUTE_KEYS.contains( + || !INCREMENTAL_TOPIC_SUPPORTED_ATTRIBUTE_KEYS.contains( key.trim().toLowerCase(Locale.ROOT))) .map(String::valueOf) .sorted() @@ -416,15 +416,15 @@ private void validateConsensusTopicAttributes(final TopicConfig topicConfig) final String exceptionMessage = String.format( ConfigNodeMessages - .EXCEPTION_FAILED_TO_CREATE_OR_ALTER_TOPIC_MODE_CONSENSUS_DOES_NOT_SUPPORT_TOPIC_ATTRIBUTES_ARG_3C2D0BDA, + .EXCEPTION_FAILED_TO_CREATE_OR_ALTER_TOPIC_MODE_INCREMENTAL_DOES_NOT_SUPPORT_TOPIC_ATTRIBUTES_ARG_1A72326A, unsupportedAttributes); LOGGER.warn(exceptionMessage); throw new SubscriptionException(exceptionMessage); } - private void validateConsensusProtocolSupport(final TopicConfig topicConfig) + private void validateIncrementalProtocolSupport(final TopicConfig topicConfig) throws SubscriptionException { - if (!topicConfig.isConsensusMode()) { + if (!topicConfig.isIncrementalMode()) { return; } @@ -437,7 +437,7 @@ private void validateConsensusProtocolSupport(final TopicConfig topicConfig) String.format( "Failed to create or alter topic, %s=%s is only supported when %s=%s, but current value is %s", TopicConstant.MODE_KEY, - TopicConstant.MODE_CONSENSUS_VALUE, + TopicConstant.MODE_INCREMENTAL_VALUE, DATA_REGION_CONSENSUS_PROTOCOL_CLASS_KEY, ConsensusFactory.IOT_CONSENSUS, actualProtocol); @@ -491,22 +491,24 @@ private void validateColumnFilter(final TopicConfig topicConfig) throws Subscrip } } - private boolean isConsensusBasedTopicConfig(final TopicConfig topicConfig) { - return topicConfig.isConsensusMode(); + private boolean isIncrementalTopicConfig(final TopicConfig topicConfig) { + return topicConfig.isIncrementalMode(); } - private void validateConsensusTopicRetentionConfig(final TopicConfig topicConfig) + private void validateIncrementalTopicRetentionConfig(final TopicConfig topicConfig) throws SubscriptionException { if (!topicConfig.hasAttribute(TopicConstant.RETENTION_BYTES_KEY) && !topicConfig.hasAttribute(TopicConstant.RETENTION_MS_KEY)) { return; } - if (!isConsensusBasedTopicConfig(topicConfig)) { + if (!isIncrementalTopicConfig(topicConfig)) { final String exceptionMessage = String.format( - "Failed to create or alter topic, %s and %s are only supported for consensus topics", - TopicConstant.RETENTION_BYTES_KEY, TopicConstant.RETENTION_MS_KEY); + ConfigNodeMessages + .EXCEPTION_FAILED_TO_CREATE_OR_ALTER_TOPIC_ARG_AND_ARG_ARE_ONLY_SUPPORTED_FOR_INCREMENTAL_TOPICS_D86CEA8E, + TopicConstant.RETENTION_BYTES_KEY, + TopicConstant.RETENTION_MS_KEY); LOGGER.warn(exceptionMessage); throw new SubscriptionException(exceptionMessage); } diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/runtime/SubscriptionHandleLeaderChangeProcedure.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/runtime/SubscriptionHandleLeaderChangeProcedure.java index 9c0801b0f86d1..c1af7445de314 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/runtime/SubscriptionHandleLeaderChangeProcedure.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/runtime/SubscriptionHandleLeaderChangeProcedure.java @@ -94,7 +94,7 @@ public boolean executeFromValidate(final ConfigNodeProcedureEnv env) { return false; } for (final TopicMeta topicMeta : subscriptionInfo.get().getAllTopicMeta()) { - if (topicMeta.getConfig().isConsensusMode()) { + if (topicMeta.getConfig().isIncrementalMode()) { return true; } } diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/CreateSubscriptionProcedure.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/CreateSubscriptionProcedure.java index 5a3904986ba38..a81661bff2fdc 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/CreateSubscriptionProcedure.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/CreateSubscriptionProcedure.java @@ -118,7 +118,7 @@ protected boolean executeFromValidate(final ConfigNodeProcedureEnv env) final TopicMeta topicMeta = subscriptionInfo.get().deepCopyTopicMeta(topicName); final String topicMode = topicMeta.getConfig().getMode(); - final boolean isConsensusBasedTopic = topicMeta.getConfig().isConsensusMode(); + final boolean isConsensusBasedTopic = topicMeta.getConfig().isIncrementalMode(); if (isConsensusBasedTopic) { // skip pipe creation diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/DropSubscriptionProcedure.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/DropSubscriptionProcedure.java index 2329a81d83dd2..321ebc9bb966e 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/DropSubscriptionProcedure.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/subscription/subscription/DropSubscriptionProcedure.java @@ -108,7 +108,7 @@ protected boolean executeFromValidate(final ConfigNodeProcedureEnv env) if (topicsUnsubByGroup.contains(topic)) { final TopicMeta topicMeta = subscriptionInfo.get().deepCopyTopicMeta(topic); final String topicMode = topicMeta.getConfig().getMode(); - final boolean isConsensusBasedTopic = topicMeta.getConfig().isConsensusMode(); + final boolean isConsensusBasedTopic = topicMeta.getConfig().isIncrementalMode(); if (isConsensusBasedTopic) { LOGGER.info( diff --git a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/persistence/subscription/SubscriptionInfoTopicValidationTest.java b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/persistence/subscription/SubscriptionInfoTopicValidationTest.java index 858bc41d75835..bda9cf81d7f2d 100644 --- a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/persistence/subscription/SubscriptionInfoTopicValidationTest.java +++ b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/persistence/subscription/SubscriptionInfoTopicValidationTest.java @@ -38,7 +38,7 @@ public class SubscriptionInfoTopicValidationTest { @Test public void testValidateColumnFilterOnCreate() throws Exception { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map attributes = newConsensusTableTopicAttributes(); + final Map attributes = newIncrementalTableTopicAttributes(); attributes.put(TopicConstant.COLUMN_FILTER_KEY, "column_name IN (\"id1\", \"m1\")"); Assert.assertTrue( @@ -58,7 +58,7 @@ public void testRejectColumnFilterOnTreeTopic() { @Test public void testColumnFilterKeyIsCaseInsensitiveOnCreate() throws Exception { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map attributes = newLiveTableTopicAttributes(); + final Map attributes = newInitialTableTopicAttributes(); attributes.put("Column-Filter", "column_name = \"id1\""); Assert.assertTrue( @@ -69,7 +69,7 @@ public void testColumnFilterKeyIsCaseInsensitiveOnCreate() throws Exception { @Test public void testRejectDuplicateColumnFilterKeys() { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map attributes = newLiveTableTopicAttributes(); + final Map attributes = newInitialTableTopicAttributes(); attributes.put(TopicConstant.COLUMN_FILTER_KEY, "column_name = \"id1\""); attributes.put("Column-Filter", "column_name = \"m1\""); @@ -88,16 +88,16 @@ public void testRejectMixedCaseColumnFilterOnTreeTopic() { @Test public void testRejectDuplicateTopicConfigKeys() { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map attributes = newLiveTableTopicAttributes(); + final Map attributes = newInitialTableTopicAttributes(); attributes.put("Mode", TopicConstant.MODE_SNAPSHOT_VALUE); assertCreateRejected(subscriptionInfo, attributes, "duplicate mode"); } @Test - public void testAcceptColumnFilterOnLiveTsFileTableTopic() throws Exception { + public void testAcceptColumnFilterOnInitialTsFileTableTopic() throws Exception { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map attributes = newLiveTableTopicAttributes(); + final Map attributes = newInitialTableTopicAttributes(); attributes.put(TopicConstant.FORMAT_KEY, TopicConstant.FORMAT_TS_FILE_VALUE); attributes.put(TopicConstant.COLUMN_FILTER_KEY, "column_name = \"id1\""); @@ -107,18 +107,18 @@ public void testAcceptColumnFilterOnLiveTsFileTableTopic() throws Exception { } @Test - public void testRejectLegacyTsFileAliasOnConsensusTopic() { + public void testRejectLegacyTsFileAliasOnIncrementalTopic() { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map attributes = newConsensusTableTopicAttributes(); + final Map attributes = newIncrementalTableTopicAttributes(); attributes.put(TopicConstant.FORMAT_KEY, "TsFileHandler"); - assertCreateRejected(subscriptionInfo, attributes, "mode=consensus only supports format"); + assertCreateRejected(subscriptionInfo, attributes, "mode=incremental only supports format"); } @Test - public void testRejectUnsupportedAttributesOnConsensusTopic() { + public void testRejectUnsupportedAttributesOnIncrementalTopic() { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map attributes = newConsensusTableTopicAttributes(); + final Map attributes = newIncrementalTableTopicAttributes(); attributes.put(TopicConstant.START_TIME_KEY, "0"); attributes.put(TopicConstant.STRICT_KEY, "false"); attributes.put("processor", "custom-processor"); @@ -126,25 +126,25 @@ public void testRejectUnsupportedAttributesOnConsensusTopic() { assertCreateRejected( subscriptionInfo, attributes, - "mode=consensus does not support topic attributes [processor, start-time, strict]"); + "mode=incremental does not support topic attributes [processor, start-time, strict]"); } @Test - public void testRejectUnknownAttributeOnConsensusTopic() { + public void testRejectUnknownAttributeOnIncrementalTopic() { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map attributes = newConsensusTableTopicAttributes(); + final Map attributes = newIncrementalTableTopicAttributes(); attributes.put("unknown-attribute", "value"); assertCreateRejected( subscriptionInfo, attributes, - "mode=consensus does not support topic attributes [unknown-attribute]"); + "mode=incremental does not support topic attributes [unknown-attribute]"); } @Test - public void testAllowPipeAttributesOnLiveTopic() throws Exception { + public void testAllowPipeAttributesOnInitialTopic() throws Exception { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map attributes = newLiveTableTopicAttributes(); + final Map attributes = newInitialTableTopicAttributes(); attributes.put(TopicConstant.START_TIME_KEY, "0"); attributes.put(TopicConstant.STRICT_KEY, "false"); attributes.put("processor", "custom-processor"); @@ -157,7 +157,7 @@ public void testAllowPipeAttributesOnLiveTopic() throws Exception { @Test public void testRejectEmptyColumnFilter() { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map attributes = newConsensusTableTopicAttributes(); + final Map attributes = newIncrementalTableTopicAttributes(); attributes.put(TopicConstant.COLUMN_FILTER_KEY, " "); assertCreateRejected(subscriptionInfo, attributes, "column-filter should not be empty"); @@ -166,12 +166,12 @@ public void testRejectEmptyColumnFilter() { @Test public void testAcceptAlteringColumnFilter() throws Exception { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map originalAttributes = newConsensusTableTopicAttributes(); + final Map originalAttributes = newIncrementalTableTopicAttributes(); originalAttributes.put(TopicConstant.COLUMN_FILTER_KEY, "column_name = \"id1\""); subscriptionInfo.createTopic( new CreateTopicPlan(new TopicMeta("table_topic", 1L, originalAttributes))); - final Map updatedAttributes = newConsensusTableTopicAttributes(); + final Map updatedAttributes = newIncrementalTableTopicAttributes(); updatedAttributes.put(TopicConstant.COLUMN_FILTER_KEY, "column_name = \"m1\""); subscriptionInfo.validateBeforeAlteringTopic( @@ -181,7 +181,7 @@ public void testAcceptAlteringColumnFilter() throws Exception { @Test public void testValidateRetentionConfigOnCreate() throws Exception { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map attributes = newConsensusTableTopicAttributes(); + final Map attributes = newIncrementalTableTopicAttributes(); attributes.put(TopicConstant.RETENTION_BYTES_KEY, "1048576"); attributes.put(TopicConstant.RETENTION_MS_KEY, "-1"); @@ -193,17 +193,17 @@ public void testValidateRetentionConfigOnCreate() throws Exception { @Test public void testRejectRetentionOnTsFileTopic() { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map attributes = newConsensusTableTopicAttributes(); + final Map attributes = newIncrementalTableTopicAttributes(); attributes.put(TopicConstant.FORMAT_KEY, TopicConstant.FORMAT_TS_FILE_VALUE); attributes.put(TopicConstant.RETENTION_BYTES_KEY, "1024"); - assertCreateRejected(subscriptionInfo, attributes, "mode=consensus only supports format"); + assertCreateRejected(subscriptionInfo, attributes, "mode=incremental only supports format"); } @Test public void testRejectIllegalRetentionValue() { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map attributes = newConsensusTableTopicAttributes(); + final Map attributes = newIncrementalTableTopicAttributes(); attributes.put(TopicConstant.RETENTION_BYTES_KEY, "0"); assertCreateRejected(subscriptionInfo, attributes, "expected -1 or a positive long value"); @@ -212,7 +212,7 @@ public void testRejectIllegalRetentionValue() { @Test public void testRejectIllegalRetentionFormat() { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map attributes = newConsensusTableTopicAttributes(); + final Map attributes = newIncrementalTableTopicAttributes(); attributes.put(TopicConstant.RETENTION_MS_KEY, "1h"); assertCreateRejected(subscriptionInfo, attributes, "expected a long value"); @@ -221,12 +221,12 @@ public void testRejectIllegalRetentionFormat() { @Test public void testRejectAlteringRetentionConfig() throws Exception { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map originalAttributes = newConsensusTableTopicAttributes(); + final Map originalAttributes = newIncrementalTableTopicAttributes(); originalAttributes.put(TopicConstant.RETENTION_BYTES_KEY, "1024"); subscriptionInfo.createTopic( new CreateTopicPlan(new TopicMeta("table_topic", 1L, originalAttributes))); - final Map updatedAttributes = newConsensusTableTopicAttributes(); + final Map updatedAttributes = newIncrementalTableTopicAttributes(); updatedAttributes.put(TopicConstant.RETENTION_BYTES_KEY, "2048"); try { @@ -247,10 +247,41 @@ public void testRejectIllegalMode() { assertCreateRejected(subscriptionInfo, attributes, "unsupported mode"); } + @SuppressWarnings("deprecation") @Test - public void testAcceptColumnFilterOnLiveTableTopic() throws Exception { + public void testAcceptLegacyModeValues() throws Exception { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map attributes = newLiveTableTopicAttributes(); + + final Map liveAttributes = newInitialTableTopicAttributes(); + liveAttributes.put(TopicConstant.MODE_KEY, TopicConstant.MODE_LIVE_VALUE); + Assert.assertTrue( + subscriptionInfo.validateBeforeCreatingTopic( + new TCreateTopicReq("live_topic").setTopicAttributes(liveAttributes))); + + final Map consensusAttributes = newIncrementalTableTopicAttributes(); + consensusAttributes.put(TopicConstant.MODE_KEY, TopicConstant.MODE_CONSENSUS_VALUE); + Assert.assertTrue( + subscriptionInfo.validateBeforeCreatingTopic( + new TCreateTopicReq("consensus_topic").setTopicAttributes(consensusAttributes))); + } + + @SuppressWarnings("deprecation") + @Test + public void testAllowAlteringModeFromLegacyAliasToCanonicalValue() throws Exception { + final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); + final Map originalAttributes = newInitialTableTopicAttributes(); + originalAttributes.put(TopicConstant.MODE_KEY, TopicConstant.MODE_LIVE_VALUE); + subscriptionInfo.createTopic( + new CreateTopicPlan(new TopicMeta("table_topic", 1L, originalAttributes))); + + subscriptionInfo.validateBeforeAlteringTopic( + new TopicMeta("table_topic", 2L, newInitialTableTopicAttributes())); + } + + @Test + public void testAcceptColumnFilterOnInitialTableTopic() throws Exception { + final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); + final Map attributes = newInitialTableTopicAttributes(); attributes.put(TopicConstant.COLUMN_FILTER_KEY, "column_name = \"id1\""); Assert.assertTrue( @@ -259,12 +290,12 @@ public void testAcceptColumnFilterOnLiveTableTopic() throws Exception { } @Test - public void testRejectConsensusOnlyRetentionOnLiveTopic() { + public void testRejectIncrementalOnlyRetentionOnInitialTopic() { final SubscriptionInfo subscriptionInfo = new SubscriptionInfo(); - final Map attributes = newLiveTableTopicAttributes(); + final Map attributes = newInitialTableTopicAttributes(); attributes.put(TopicConstant.RETENTION_BYTES_KEY, "1024"); - assertCreateRejected(subscriptionInfo, attributes, "only supported for consensus topics"); + assertCreateRejected(subscriptionInfo, attributes, "only supported for incremental topics"); } @Test @@ -294,18 +325,18 @@ public void testAcceptOwnerLeaseDurationAtMin() throws Exception { new TCreateTopicReq("owner_topic").setTopicAttributes(attributes))); } - private static Map newConsensusTableTopicAttributes() { + private static Map newIncrementalTableTopicAttributes() { final Map attributes = new HashMap<>(); attributes.put(SystemConstant.SQL_DIALECT_KEY, SystemConstant.SQL_DIALECT_TABLE_VALUE); - attributes.put(TopicConstant.MODE_KEY, TopicConstant.MODE_CONSENSUS_VALUE); + attributes.put(TopicConstant.MODE_KEY, TopicConstant.MODE_INCREMENTAL_VALUE); attributes.put(TopicConstant.FORMAT_KEY, TopicConstant.FORMAT_RECORD_HANDLER_VALUE); return attributes; } - private static Map newLiveTableTopicAttributes() { + private static Map newInitialTableTopicAttributes() { final Map attributes = new HashMap<>(); attributes.put(SystemConstant.SQL_DIALECT_KEY, SystemConstant.SQL_DIALECT_TABLE_VALUE); - attributes.put(TopicConstant.MODE_KEY, TopicConstant.MODE_LIVE_VALUE); + attributes.put(TopicConstant.MODE_KEY, TopicConstant.MODE_INITIAL_VALUE); attributes.put(TopicConstant.FORMAT_KEY, TopicConstant.FORMAT_RECORD_HANDLER_VALUE); return attributes; } diff --git a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java index 8e61bcda388c3..757529f3aa0c4 100644 --- a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java +++ b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java @@ -2020,10 +2020,15 @@ private DataNodePipeMessages() {} + "runtimeVersion {} -> {}, runtimeState={} (route hint)"; public static final String PIPE_LOG_FAILED_TO_CHECK_IF_TOPIC_IS_CONSENSUS_BASED_DEFAULTING_TO_ECCE1509 = "Failed to check if topic [{}] is consensus-based, defaulting to false"; - public static final String PIPE_LOG_SKIPPING_SETUP_OF_CONSENSUS_BASED_SUBSCRIPTIONS_FOR_CONSUMER_A7B2C812 = + public static final String PIPE_LOG_SKIPPING_SETUP_OF_CONSENSUS_BASED_SUBSCRIPTIONS_FOR_CONSUMER_46BEE6E4 = "Skipping setup of consensus-based subscriptions for consumer group [{}] because " - + "mode=consensus only supports data_region_consensus_protocol_class={}, but current " + + "mode=incremental only supports data_region_consensus_protocol_class={}, but current " + "configured value is {} (runtime consensus implementation: {})"; + public static final String + EXCEPTION_SUBSCRIPTION_CANNOT_ARG_CONSENSUS_BASED_TOPIC_S_ARG_IN_CONSUMER_GROUP_ARG_BECAUSE_MODE_INCREMENTAL_ONLY_SUPPORTS_DATA_REGION_CONSENSUS_PROTOCOL_CLASS_ARG_BUT_CURRENT_CONFIGURED_VALUE_IS_ARG_RUNTIME_CONSENSUS_IMPLEMENTATION_ARG_6F21ED67 = + "Subscription: cannot %s consensus-based topic(s) %s in consumer group [%s] because " + + "mode=incremental only supports data_region_consensus_protocol_class=%s, but " + + "current configured value is %s (runtime consensus implementation: %s)"; public static final String PIPE_LOG_TOPIC_CONFIG_NOT_FOUND_FOR_TOPIC_CANNOT_SET_UP_CONSENSUS_A93339CE = "Topic config not found for topic [{}], cannot set up consensus queue"; public static final String PIPE_LOG_NO_LOCAL_IOTCONSENSUS_DATA_REGION_FOUND_FOR_TOPIC_IN_CONSUMER_6FD0600E = diff --git a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java index 2389f97c8ec01..6fd1922c6a6fb 100644 --- a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java +++ b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java @@ -1875,9 +1875,14 @@ private DataNodePipeMessages() {} + "{} -> {},runtimeState={}(route hint)"; public static final String PIPE_LOG_FAILED_TO_CHECK_IF_TOPIC_IS_CONSENSUS_BASED_DEFAULTING_TO_ECCE1509 = "检查 topic [{}] 是否为 consensus-based 失败,默认设为 false"; - public static final String PIPE_LOG_SKIPPING_SETUP_OF_CONSENSUS_BASED_SUBSCRIPTIONS_FOR_CONSUMER_A7B2C812 = - "跳过 consumer group [{}] 的 consensus-based subscription 设置,因为 mode=consensus 仅支持 " + public static final String PIPE_LOG_SKIPPING_SETUP_OF_CONSENSUS_BASED_SUBSCRIPTIONS_FOR_CONSUMER_46BEE6E4 = + "跳过 consumer group [{}] 的 consensus-based subscription 设置,因为 mode=incremental 仅支持 " + "data_region_consensus_protocol_class={},但当前配置值为 {}(运行时 consensus 实现:{})"; + public static final String + EXCEPTION_SUBSCRIPTION_CANNOT_ARG_CONSENSUS_BASED_TOPIC_S_ARG_IN_CONSUMER_GROUP_ARG_BECAUSE_MODE_INCREMENTAL_ONLY_SUPPORTS_DATA_REGION_CONSENSUS_PROTOCOL_CLASS_ARG_BUT_CURRENT_CONFIGURED_VALUE_IS_ARG_RUNTIME_CONSENSUS_IMPLEMENTATION_ARG_6F21ED67 = + "Subscription:无法执行 %s,consensus-based topic 为 %s,consumer group 为 [%s],因为 " + + "mode=incremental 仅支持 data_region_consensus_protocol_class=%s,但当前配置值为 %s" + + "(运行时 consensus 实现:%s)"; public static final String PIPE_LOG_TOPIC_CONFIG_NOT_FOUND_FOR_TOPIC_CANNOT_SET_UP_CONSENSUS_A93339CE = "未找到 topic [{}] 的配置,无法设置 consensus queue"; public static final String PIPE_LOG_NO_LOCAL_IOTCONSENSUS_DATA_REGION_FOUND_FOR_TOPIC_IN_CONSUMER_6FD0600E = diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/executor/ClusterConfigTaskExecutor.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/executor/ClusterConfigTaskExecutor.java index baef673ca8605..65d4b3f102d5e 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/executor/ClusterConfigTaskExecutor.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/executor/ClusterConfigTaskExecutor.java @@ -3172,7 +3172,7 @@ public SettableFuture createTopic( // Validate topic config final TopicMeta temporaryTopicMeta = new TopicMeta(topicName, System.currentTimeMillis(), topicAttributes); - if (!temporaryTopicMeta.getConfig().isConsensusMode()) { + if (!temporaryTopicMeta.getConfig().isIncrementalMode()) { try { PipeDataNodeAgent.plugin() .validate( diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionBrokerAgent.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionBrokerAgent.java index 2fb64260e502d..019253774d5bb 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionBrokerAgent.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/agent/SubscriptionBrokerAgent.java @@ -422,9 +422,8 @@ private String buildUnsupportedConsensusRuntimeMessage( final String runtimeConsensusImplementation = Objects.nonNull(dataRegionConsensus) ? dataRegionConsensus.getClass().getName() : "null"; return String.format( - "Subscription: cannot %s consensus-based topic(s) %s in consumer group [%s] because " - + "mode=consensus only supports data_region_consensus_protocol_class=%s, but current " - + "configured value is %s (runtime consensus implementation: %s)", + DataNodePipeMessages + .EXCEPTION_SUBSCRIPTION_CANNOT_ARG_CONSENSUS_BASED_TOPIC_S_ARG_IN_CONSUMER_GROUP_ARG_BECAUSE_MODE_INCREMENTAL_ONLY_SUPPORTS_DATA_REGION_CONSENSUS_PROTOCOL_CLASS_ARG_BUT_CURRENT_CONFIGURED_VALUE_IS_ARG_RUNTIME_CONSENSUS_IMPLEMENTATION_ARG_6F21ED67, operation, topicNames, consumerGroupId, diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusSubscriptionSetupHandler.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusSubscriptionSetupHandler.java index 390d5fd9f2e72..21fbd31724115 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusSubscriptionSetupHandler.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusSubscriptionSetupHandler.java @@ -61,9 +61,9 @@ /** * Handles setup and teardown of consensus-based subscription queues on DataNode. * - *

For each consensus-mode topic subscribed by a consumer group, this handler discovers matching - * local IoTConsensus DataRegions, builds the appropriate log-to-tablet converter, and binds one - * queue per region to the consensus subscription broker. + *

For each incremental-mode topic subscribed by a consumer group, this handler discovers + * matching local IoTConsensus DataRegions, builds the appropriate log-to-tablet converter, and + * binds one queue per region to the consensus subscription broker. */ public class ConsensusSubscriptionSetupHandler { @@ -294,7 +294,7 @@ private static void onRegionRemoved(final ConsensusGroupId groupId) { public static boolean isConsensusBasedTopic(final String topicName) { try { final String topicMode = SubscriptionAgent.topic().getTopicMode(topicName); - final boolean result = TopicConstant.MODE_CONSENSUS_VALUE.equalsIgnoreCase(topicMode); + final boolean result = TopicConstant.MODE_INCREMENTAL_VALUE.equalsIgnoreCase(topicMode); LOGGER.debug( DataNodePipeMessages.PIPE_LOG_ISCONSENSUSBASEDTOPIC_CHECK_FOR_TOPIC_MODE_RESULT_19EFA0F9, topicName, @@ -320,7 +320,7 @@ private static boolean isConsensusBasedTopicRequired(final String topicName) { .EXCEPTION_TOPIC_METADATA_FOR_ARG_IS_UNAVAILABLE_DURING_CONSENSUS_SUBSCRIPTION_SETUP_A1949F20, topicName)); } - return TopicConstant.MODE_CONSENSUS_VALUE.equalsIgnoreCase(topicMode); + return TopicConstant.MODE_INCREMENTAL_VALUE.equalsIgnoreCase(topicMode); } public static void setupConsensusSubscriptions( @@ -332,7 +332,7 @@ public static void setupConsensusSubscriptions( Objects.nonNull(dataRegionConsensus) ? dataRegionConsensus.getClass().getName() : "null"; LOGGER.warn( DataNodePipeMessages - .PIPE_LOG_SKIPPING_SETUP_OF_CONSENSUS_BASED_SUBSCRIPTIONS_FOR_CONSUMER_A7B2C812, + .PIPE_LOG_SKIPPING_SETUP_OF_CONSENSUS_BASED_SUBSCRIPTIONS_FOR_CONSUMER_46BEE6E4, consumerGroupId, ConsensusFactory.IOT_CONSENSUS, configuredProtocol,