This is an automated email from the ASF dual-hosted git repository.
jt2594838 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new d82806fa208 Subscription: Rename topic modes to initial and
incremental (#18370)
d82806fa208 is described below
commit d82806fa208cc0745271e1540c8100f7e4a3f54b
Author: Caideyipi <[email protected]>
AuthorDate: Mon Aug 3 11:41:25 2026 +0800
Subscription: Rename topic modes to initial and incremental (#18370)
---
.../iotdb/ConsensusSubscriptionSessionExample.java | 2 +-
...sensusTableModelSubscriptionSessionExample.java | 2 +-
.../local/ConsensusSubscriptionITSupport.java | 2 +-
.../ConsensusSubscriptionTableITSupport.java | 2 +-
...ConsensusSubscriptionColumnFilterClusterIT.java | 2 +-
.../IoTDBSubscriptionColumnFilterIT.java | 5 +-
.../tablemodel/IoTDBSubscriptionPermissionIT.java | 2 +-
.../IoTDBAllTsDatasetPullConsumerIT.java | 2 +-
.../IoTDBAllTsTsfilePullConsumerIT.java | 2 +-
.../IoTDBPathDeviceDataSetPullConsumerIT.java | 2 +-
.../IoTDBPathDeviceTsfilePullConsumerIT.java | 2 +-
.../IoTDBTimeTsDatasetPullConsumerIT.java | 2 +-
.../IoTDBTimeTsTsfilePullConsumerIT.java | 2 +-
.../IoTDBLooseAllTsDatasetPushConsumerIT.java | 4 +-
.../IoTDBLooseAllTsfilePushConsumerIT.java | 4 +-
.../IoTDBPathLooseDeviceTsfilePushConsumerIT.java | 2 +-
.../IoTDBPathLooseTsDatasetPushConsumerIT.java | 2 +-
.../IoTDBPathLooseTsfilePushConsumerIT.java | 2 +-
.../IoTDBPathTsLooseDatasetPushConsumerIT.java | 2 +-
.../IoTDBTimeLooseTsDatasetPushConsumerIT.java | 4 +-
.../IoTDBTimeLooseTsTsfilePushConsumerIT.java | 4 +-
.../IoTDBTimeLooseTsfilePushConsumerIT.java | 2 +-
.../IoTDBTimeTsLooseDatasetPushConsumerIT.java | 4 +-
.../subscription/i18n/SubscriptionMessages.java | 2 +-
.../subscription/i18n/SubscriptionMessages.java | 2 +-
.../iotdb/rpc/subscription/config/TopicConfig.java | 51 ++++++++--
.../rpc/subscription/config/TopicConstant.java | 19 +++-
.../base/AbstractSubscriptionConsumer.java | 2 +-
.../rpc/subscription/config/TopicConfigTest.java | 57 ++++++++++++
.../iotdb/confignode/i18n/ConfigNodeMessages.java | 7 +-
.../iotdb/confignode/i18n/ConfigNodeMessages.java | 7 +-
.../runtime/SubscriptionRuntimeCoordinator.java | 2 +-
.../persistence/subscription/SubscriptionInfo.java | 44 ++++-----
.../SubscriptionHandleLeaderChangeProcedure.java | 2 +-
.../subscription/CreateSubscriptionProcedure.java | 2 +-
.../subscription/DropSubscriptionProcedure.java | 2 +-
.../SubscriptionInfoTopicValidationTest.java | 103 ++++++++++++++-------
.../apache/iotdb/db/i18n/DataNodePipeMessages.java | 9 +-
.../apache/iotdb/db/i18n/DataNodePipeMessages.java | 9 +-
.../config/executor/ClusterConfigTaskExecutor.java | 2 +-
.../agent/SubscriptionBrokerAgent.java | 5 +-
.../ConsensusSubscriptionSetupHandler.java | 12 +--
42 files changed, 274 insertions(+), 124 deletions(-)
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 c0ebbe37198..12596f5d5b9 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 @@ public class ConsensusSubscriptionSessionExample {
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 a877a4a861e..fa7733d57f4 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 @@ public class ConsensusTableModelSubscriptionSessionExample {
.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 254b5ffeb85..066c38ef831 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 @@ final class ConsensusSubscriptionITSupport {
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 4bd38992ed1..29a859d3011 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 @@ final class ConsensusSubscriptionTableITSupport {
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 d2fd5552ee7..6c33532828b 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 @@ public class
IoTDBConsensusSubscriptionColumnFilterClusterIT extends AbstractSub
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 27cd10ed534..bfeb7e6b7aa 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 @@ public class IoTDBSubscriptionColumnFilterIT extends
AbstractSubscriptionDualIT
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 @@ public class IoTDBSubscriptionColumnFilterIT extends
AbstractSubscriptionDualIT
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 123a68c8df5..44d0dea16d7 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 class IoTDBSubscriptionPermissionIT extends
AbstractSubscriptionLocalIT {
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 a1a0cae5d9c..fe2ac71aba7 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 class IoTDBAllTsDatasetPullConsumerIT extends
AbstractSubscriptionTreeReg
"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 c00ea2401ce..8a585c5727e 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 class IoTDBAllTsTsfilePullConsumerIT extends
AbstractSubscriptionTreeRegr
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 2cb4b757253..a048041ce90 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 class IoTDBPathDeviceDataSetPullConsumerIT extends
AbstractSubscriptionTr
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 70a5857ac88..746e2abdcb4 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 class IoTDBPathDeviceTsfilePullConsumerIT extends
AbstractSubscriptionTre
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 3fc35e39bb9..1b8e9174e9b 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 class IoTDBTimeTsDatasetPullConsumerIT extends
AbstractSubscriptionTreeRe
"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 7cec98a6270..f609702e431 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 class IoTDBTimeTsTsfilePullConsumerIT extends
AbstractSubscriptionTreeReg
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 c30e9d8493d..af8d620637f 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 @@ import static
org.apache.iotdb.subscription.it.IoTDBSubscriptionITConstant.AWAIT
* DataSet
* pattern: ts
* loose-range: all
- * mode: live
+ * mode: initial
*/
@RunWith(IoTDBTestRunner.class)
@Category({MultiClusterIT2SubscriptionTreeRegressionConsumer.class})
@@ -82,7 +82,7 @@ public class IoTDBLooseAllTsDatasetPushConsumerIT extends
AbstractSubscriptionTr
"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 91fb1a281e5..1c2f420f339 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 @@ import static
org.apache.iotdb.subscription.it.IoTDBSubscriptionITConstant.AWAIT
/***
* push consumer
- * mode: live
+ * mode: initial
* pattern: db
* loose-range: all
*/
@@ -84,7 +84,7 @@ public class IoTDBLooseAllTsfilePushConsumerIT extends
AbstractSubscriptionTreeR
"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 c647160493a..270d5cd928f 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 class IoTDBPathLooseDeviceTsfilePushConsumerIT extends
AbstractSubscripti
"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 3e144c137fd..de854f29317 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 class IoTDBPathLooseTsDatasetPushConsumerIT extends
AbstractSubscriptionT
"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 b8c438bfb9d..c8800aa93ac 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 class IoTDBPathLooseTsfilePushConsumerIT extends
AbstractSubscriptionTree
"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 5a2a5d6a351..aba24718a63 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 class IoTDBPathTsLooseDatasetPushConsumerIT extends
AbstractSubscriptionT
"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 ecbd661e5e0..38b826f3caf 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 @@ import static
org.apache.iotdb.subscription.it.IoTDBSubscriptionITConstant.AWAIT
* DataSet
* pattern: ts
* loose-range: time
- * live
+ * initial
*/
@RunWith(IoTDBTestRunner.class)
@Category({MultiClusterIT2SubscriptionTreeRegressionConsumer.class})
@@ -84,7 +84,7 @@ public class IoTDBTimeLooseTsDatasetPushConsumerIT extends
AbstractSubscriptionT
"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 0311bc0a0b5..e62644a59dc 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 java.util.concurrent.atomic.AtomicInteger;
import static
org.apache.iotdb.subscription.it.IoTDBSubscriptionITConstant.AWAIT;
/***
- * mode: live
+ * mode: initial
* loose-range:path
* format: tsfile
*/
@@ -83,7 +83,7 @@ public class IoTDBTimeLooseTsTsfilePushConsumerIT extends
AbstractSubscriptionTr
"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 53207ee7fe6..68b36104b01 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 class IoTDBTimeLooseTsfilePushConsumerIT extends
AbstractSubscriptionTree
"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 4f358ddb8c4..bfd539054e9 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 @@ import static
org.apache.iotdb.subscription.it.IoTDBSubscriptionITConstant.AWAIT
* DataSet
* pattern: ts
* time loose
- * live
+ * initial
*/
@RunWith(IoTDBTestRunner.class)
@Category({MultiClusterIT2SubscriptionTreeRegressionConsumer.class})
@@ -83,7 +83,7 @@ public class IoTDBTimeTsLooseDatasetPushConsumerIT extends
AbstractSubscriptionT
"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 7ce7e7a8763..95f26d55d68 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 @@ public final class 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 4a8edd3f54a..562a4d72572 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 @@ public final class 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 b0824fcb394..3d107476b85 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.nio.ByteBuffer;
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<String> 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<String> orderModes = new HashSet<>(3);
@@ -84,8 +85,10 @@ public class TopicConfig extends PipeParameters {
private static final Map<String, String> SNAPSHOT_MODE_CONFIG =
Collections.singletonMap("mode", TopicConstant.MODE_SNAPSHOT_VALUE);
- private static final Map<String, String> 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<String, String> INITIAL_MODE_CONFIG =
+ Collections.singletonMap("mode", TopicConstant.LEGACY_MODE_LIVE_VALUE);
private static final Map<String, String> STRICT_MODE_CONFIG =
Collections.singletonMap("mode.strict", "true");
@@ -125,12 +128,28 @@ public class TopicConfig extends PipeParameters {
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 class TopicConfig extends PipeParameters {
}
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 class TopicConfig extends PipeParameters {
}
public Map<String, String> 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<String, String> 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 f929d7b7472..4d11f2fe8a2 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 74491898e59..fa48b6a099a 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 @@ abstract class AbstractSubscriptionConsumer implements
AutoCloseable {
private boolean allTopicMessagesHaveBeenConsumed(final Collection<String>
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 d60f4f1f702..d308ef123cf 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 @@ import java.util.Map;
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 class TopicConfigTest {
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 014ad8e807c..54e06e9a972 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 @@ public final class 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 3c3bce0b845..0c28efc4096 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 @@ public final class 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 1f3a70073c5..dd8fe2e70cf 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 class SubscriptionRuntimeCoordinator {
.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 4d063e83152..12b25159fa9 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<String> CONSENSUS_TOPIC_SUPPORTED_ATTRIBUTE_KEYS =
+ private static final Set<String> INCREMENTAL_TOPIC_SUPPORTED_ATTRIBUTE_KEYS =
Set.of(
SystemConstant.SQL_DIALECT_KEY,
TopicConstant.PATH_KEY,
@@ -340,21 +340,21 @@ public class SubscriptionInfo implements
SnapshotProcessor {
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 @@ public class SubscriptionInfo implements SnapshotProcessor {
}
validateColumnFilter(topicConfig);
- validateConsensusTopicRetentionConfig(topicConfig);
+ validateIncrementalTopicRetentionConfig(topicConfig);
final Long ownerLeaseDurationMs =
topicConfig.getLong(TopicConstant.OWNER_LEASE_DURATION_MS_KEY);
@@ -393,9 +393,9 @@ public class SubscriptionInfo implements SnapshotProcessor {
}
}
- 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 @@ public class SubscriptionInfo implements SnapshotProcessor {
.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 @@ public class SubscriptionInfo implements
SnapshotProcessor {
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 @@ public class SubscriptionInfo implements SnapshotProcessor {
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 @@ public class SubscriptionInfo implements
SnapshotProcessor {
}
}
- 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 9c0801b0f86..c1af7445de3 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 class SubscriptionHandleLeaderChangeProcedure extends
AbstractOperateSubs
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 5a3904986ba..a81661bff2f 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 @@ public class CreateSubscriptionProcedure extends
AbstractOperateSubscriptionAndP
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 2329a81d83d..321ebc9bb96 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 @@ public class DropSubscriptionProcedure extends
AbstractOperateSubscriptionAndPip
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 858bc41d758..bda9cf81d7f 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<String, String> attributes = newConsensusTableTopicAttributes();
+ final Map<String, String> attributes =
newIncrementalTableTopicAttributes();
attributes.put(TopicConstant.COLUMN_FILTER_KEY, "column_name IN (\"id1\",
\"m1\")");
Assert.assertTrue(
@@ -58,7 +58,7 @@ public class SubscriptionInfoTopicValidationTest {
@Test
public void testColumnFilterKeyIsCaseInsensitiveOnCreate() throws Exception {
final SubscriptionInfo subscriptionInfo = new SubscriptionInfo();
- final Map<String, String> attributes = newLiveTableTopicAttributes();
+ final Map<String, String> attributes = newInitialTableTopicAttributes();
attributes.put("Column-Filter", "column_name = \"id1\"");
Assert.assertTrue(
@@ -69,7 +69,7 @@ public class SubscriptionInfoTopicValidationTest {
@Test
public void testRejectDuplicateColumnFilterKeys() {
final SubscriptionInfo subscriptionInfo = new SubscriptionInfo();
- final Map<String, String> attributes = newLiveTableTopicAttributes();
+ final Map<String, String> attributes = newInitialTableTopicAttributes();
attributes.put(TopicConstant.COLUMN_FILTER_KEY, "column_name = \"id1\"");
attributes.put("Column-Filter", "column_name = \"m1\"");
@@ -88,16 +88,16 @@ public class SubscriptionInfoTopicValidationTest {
@Test
public void testRejectDuplicateTopicConfigKeys() {
final SubscriptionInfo subscriptionInfo = new SubscriptionInfo();
- final Map<String, String> attributes = newLiveTableTopicAttributes();
+ final Map<String, String> 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<String, String> attributes = newLiveTableTopicAttributes();
+ final Map<String, String> 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 class SubscriptionInfoTopicValidationTest {
}
@Test
- public void testRejectLegacyTsFileAliasOnConsensusTopic() {
+ public void testRejectLegacyTsFileAliasOnIncrementalTopic() {
final SubscriptionInfo subscriptionInfo = new SubscriptionInfo();
- final Map<String, String> attributes = newConsensusTableTopicAttributes();
+ final Map<String, String> 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<String, String> attributes = newConsensusTableTopicAttributes();
+ final Map<String, String> 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 class SubscriptionInfoTopicValidationTest {
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<String, String> attributes = newConsensusTableTopicAttributes();
+ final Map<String, String> 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<String, String> attributes = newLiveTableTopicAttributes();
+ final Map<String, String> 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 class SubscriptionInfoTopicValidationTest {
@Test
public void testRejectEmptyColumnFilter() {
final SubscriptionInfo subscriptionInfo = new SubscriptionInfo();
- final Map<String, String> attributes = newConsensusTableTopicAttributes();
+ final Map<String, String> attributes =
newIncrementalTableTopicAttributes();
attributes.put(TopicConstant.COLUMN_FILTER_KEY, " ");
assertCreateRejected(subscriptionInfo, attributes, "column-filter should
not be empty");
@@ -166,12 +166,12 @@ public class SubscriptionInfoTopicValidationTest {
@Test
public void testAcceptAlteringColumnFilter() throws Exception {
final SubscriptionInfo subscriptionInfo = new SubscriptionInfo();
- final Map<String, String> originalAttributes =
newConsensusTableTopicAttributes();
+ final Map<String, String> originalAttributes =
newIncrementalTableTopicAttributes();
originalAttributes.put(TopicConstant.COLUMN_FILTER_KEY, "column_name =
\"id1\"");
subscriptionInfo.createTopic(
new CreateTopicPlan(new TopicMeta("table_topic", 1L,
originalAttributes)));
- final Map<String, String> updatedAttributes =
newConsensusTableTopicAttributes();
+ final Map<String, String> updatedAttributes =
newIncrementalTableTopicAttributes();
updatedAttributes.put(TopicConstant.COLUMN_FILTER_KEY, "column_name =
\"m1\"");
subscriptionInfo.validateBeforeAlteringTopic(
@@ -181,7 +181,7 @@ public class SubscriptionInfoTopicValidationTest {
@Test
public void testValidateRetentionConfigOnCreate() throws Exception {
final SubscriptionInfo subscriptionInfo = new SubscriptionInfo();
- final Map<String, String> attributes = newConsensusTableTopicAttributes();
+ final Map<String, String> attributes =
newIncrementalTableTopicAttributes();
attributes.put(TopicConstant.RETENTION_BYTES_KEY, "1048576");
attributes.put(TopicConstant.RETENTION_MS_KEY, "-1");
@@ -193,17 +193,17 @@ public class SubscriptionInfoTopicValidationTest {
@Test
public void testRejectRetentionOnTsFileTopic() {
final SubscriptionInfo subscriptionInfo = new SubscriptionInfo();
- final Map<String, String> attributes = newConsensusTableTopicAttributes();
+ final Map<String, String> 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<String, String> attributes = newConsensusTableTopicAttributes();
+ final Map<String, String> attributes =
newIncrementalTableTopicAttributes();
attributes.put(TopicConstant.RETENTION_BYTES_KEY, "0");
assertCreateRejected(subscriptionInfo, attributes, "expected -1 or a
positive long value");
@@ -212,7 +212,7 @@ public class SubscriptionInfoTopicValidationTest {
@Test
public void testRejectIllegalRetentionFormat() {
final SubscriptionInfo subscriptionInfo = new SubscriptionInfo();
- final Map<String, String> attributes = newConsensusTableTopicAttributes();
+ final Map<String, String> attributes =
newIncrementalTableTopicAttributes();
attributes.put(TopicConstant.RETENTION_MS_KEY, "1h");
assertCreateRejected(subscriptionInfo, attributes, "expected a long
value");
@@ -221,12 +221,12 @@ public class SubscriptionInfoTopicValidationTest {
@Test
public void testRejectAlteringRetentionConfig() throws Exception {
final SubscriptionInfo subscriptionInfo = new SubscriptionInfo();
- final Map<String, String> originalAttributes =
newConsensusTableTopicAttributes();
+ final Map<String, String> originalAttributes =
newIncrementalTableTopicAttributes();
originalAttributes.put(TopicConstant.RETENTION_BYTES_KEY, "1024");
subscriptionInfo.createTopic(
new CreateTopicPlan(new TopicMeta("table_topic", 1L,
originalAttributes)));
- final Map<String, String> updatedAttributes =
newConsensusTableTopicAttributes();
+ final Map<String, String> updatedAttributes =
newIncrementalTableTopicAttributes();
updatedAttributes.put(TopicConstant.RETENTION_BYTES_KEY, "2048");
try {
@@ -247,10 +247,41 @@ public class SubscriptionInfoTopicValidationTest {
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<String, String> attributes = newLiveTableTopicAttributes();
+
+ final Map<String, String> liveAttributes =
newInitialTableTopicAttributes();
+ liveAttributes.put(TopicConstant.MODE_KEY, TopicConstant.MODE_LIVE_VALUE);
+ Assert.assertTrue(
+ subscriptionInfo.validateBeforeCreatingTopic(
+ new
TCreateTopicReq("live_topic").setTopicAttributes(liveAttributes)));
+
+ final Map<String, String> 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<String, String> 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<String, String> attributes = newInitialTableTopicAttributes();
attributes.put(TopicConstant.COLUMN_FILTER_KEY, "column_name = \"id1\"");
Assert.assertTrue(
@@ -259,12 +290,12 @@ public class SubscriptionInfoTopicValidationTest {
}
@Test
- public void testRejectConsensusOnlyRetentionOnLiveTopic() {
+ public void testRejectIncrementalOnlyRetentionOnInitialTopic() {
final SubscriptionInfo subscriptionInfo = new SubscriptionInfo();
- final Map<String, String> attributes = newLiveTableTopicAttributes();
+ final Map<String, String> 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 class SubscriptionInfoTopicValidationTest {
new
TCreateTopicReq("owner_topic").setTopicAttributes(attributes)));
}
- private static Map<String, String> newConsensusTableTopicAttributes() {
+ private static Map<String, String> newIncrementalTableTopicAttributes() {
final Map<String, String> 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<String, String> newLiveTableTopicAttributes() {
+ private static Map<String, String> newInitialTableTopicAttributes() {
final Map<String, String> 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 8e61bcda388..757529f3aa0 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 @@ public final class 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 2389f97c8ec..6fd1922c6a6 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 @@ public final class 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 baef673ca86..65d4b3f102d 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 class ClusterConfigTaskExecutor implements
IConfigTaskExecutor {
// 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 2fb64260e50..019253774d5 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 @@ public class SubscriptionBrokerAgent {
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 390d5fd9f2e..21fbd317241 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 @@ import java.util.function.Predicate;
/**
* Handles setup and teardown of consensus-based subscription queues on
DataNode.
*
- * <p>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.
+ * <p>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 @@ public class ConsensusSubscriptionSetupHandler {
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 @@ public class ConsensusSubscriptionSetupHandler {
.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 class ConsensusSubscriptionSetupHandler {
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,