Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
992 changes: 0 additions & 992 deletions .github/workflows/daily-it.yml

Large diffs are not rendered by default.

1,102 changes: 0 additions & 1,102 deletions .github/workflows/pipe-it.yml

This file was deleted.

4 changes: 4 additions & 0 deletions integration-test/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -324,6 +324,10 @@
<configuration>
<groups>${integrationTest.includedGroups}</groups>
<excludedGroups>${integrationTest.excludedGroups}</excludedGroups>
<!-- Subscription IT is maintained and executed by the enterprise edition. -->
<excludes>
<exclude>**/subscription/it/**</exclude>
</excludes>
<useSystemClassLoader>false</useSystemClassLoader>
<parallel>none</parallel>
<forkCount>${integrationTest.forkCount}</forkCount>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -495,12 +495,6 @@ public CommonConfig setIsPipeEnableMemoryCheck(boolean isPipeEnableMemoryCheck)
return this;
}

@Override
public CommonConfig setSubscriptionEnabled(boolean subscriptionEnabled) {
setProperty("subscription_enabled", String.valueOf(subscriptionEnabled));
return this;
}

@Override
public CommonConfig setSubscriptionOwnerLeaseDurationMsMin(
long subscriptionOwnerLeaseDurationMsMin) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -512,13 +512,6 @@ public CommonConfig setIsPipeEnableMemoryCheck(boolean isPipeEnableMemoryCheck)
return this;
}

@Override
public CommonConfig setSubscriptionEnabled(boolean subscriptionEnabled) {
dnConfig.setSubscriptionEnabled(subscriptionEnabled);
cnConfig.setSubscriptionEnabled(subscriptionEnabled);
return this;
}

@Override
public CommonConfig setSubscriptionOwnerLeaseDurationMsMin(
long subscriptionOwnerLeaseDurationMsMin) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -363,11 +363,6 @@ public CommonConfig setIsPipeEnableMemoryCheck(boolean isPipeEnableMemoryCheck)
return this;
}

@Override
public CommonConfig setSubscriptionEnabled(boolean subscriptionEnabled) {
return this;
}

@Override
public CommonConfig setSubscriptionOwnerLeaseDurationMsMin(
long subscriptionOwnerLeaseDurationMsMin) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -159,8 +159,6 @@ CommonConfig setEnableAutoLeaderBalanceForRatisConsensus(

CommonConfig setIsPipeEnableMemoryCheck(boolean isPipeEnableMemoryCheck);

CommonConfig setSubscriptionEnabled(boolean subscriptionEnabled);

CommonConfig setSubscriptionOwnerLeaseDurationMsMin(long subscriptionOwnerLeaseDurationMsMin);

CommonConfig setPipeAirGapReceiverEnabled(boolean isPipeAirGapReceiverEnabled);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -358,11 +358,6 @@ public void testAsyncConnectorUseNodeUrls() throws Exception {
doTestUseNodeUrls(BuiltinPipePlugin.IOTDB_THRIFT_ASYNC_CONNECTOR.getPipePluginName());
}

@Test
public void testAirGapConnectorUseNodeUrls() throws Exception {
doTestUseNodeUrls(BuiltinPipePlugin.IOTDB_AIR_GAP_CONNECTOR.getPipePluginName());
}

private void doTestUseNodeUrls(String connectorName) throws Exception {
senderEnv
.getConfig()
Expand Down Expand Up @@ -407,16 +402,7 @@ private void doTestUseNodeUrls(String connectorName) throws Exception {
};

for (final DataNodeWrapper wrapper : receiverEnv.getDataNodeWrapperList()) {
if (connectorName.equals(BuiltinPipePlugin.IOTDB_AIR_GAP_CONNECTOR.getPipePluginName())) {
// Use default port for convenience
nodeUrlsBuilder
.append(wrapper.getIp())
.append(":")
.append(wrapper.getPipeAirGapReceiverPort())
.append(",");
} else {
nodeUrlsBuilder.append(wrapper.getIpAndPortString()).append(",");
}
nodeUrlsBuilder.append(wrapper.getIpAndPortString()).append(",");
}

try (final SyncConfigNodeIServiceClient client =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -353,11 +353,6 @@ public void testAsyncConnectorUseNodeUrls() throws Exception {
doTestUseNodeUrls(BuiltinPipePlugin.IOTDB_THRIFT_ASYNC_CONNECTOR.getPipePluginName());
}

@Test
public void testAirGapConnectorUseNodeUrls() throws Exception {
doTestUseNodeUrls(BuiltinPipePlugin.IOTDB_AIR_GAP_CONNECTOR.getPipePluginName());
}

private void doTestUseNodeUrls(String sinkName) throws Exception {
senderEnv
.getConfig()
Expand Down Expand Up @@ -398,16 +393,7 @@ private void doTestUseNodeUrls(String sinkName) throws Exception {

final StringBuilder nodeUrlsBuilder = new StringBuilder();
for (final DataNodeWrapper wrapper : receiverEnv.getDataNodeWrapperList()) {
if (sinkName.equals(BuiltinPipePlugin.IOTDB_AIR_GAP_CONNECTOR.getPipePluginName())) {
// Use default port for convenience
nodeUrlsBuilder
.append(wrapper.getIp())
.append(":")
.append(wrapper.getPipeAirGapReceiverPort())
.append(",");
} else {
nodeUrlsBuilder.append(wrapper.getIpAndPortString()).append(",");
}
nodeUrlsBuilder.append(wrapper.getIpAndPortString()).append(",");
}

try (final SyncConfigNodeIServiceClient client =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -795,8 +795,7 @@ public void testInformationSchema() throws SQLException {
"test,view_table,s11,INT32,FIELD,USING,null,",
"test,view_table,s3,INT32,FIELD,USING,null,")));

statement.execute(
"create pipe a2b with source('double-living'='true') with sink ('sink'='write-back-sink')");
statement.execute("create pipe a2b with sink ('sink'='do-nothing-sink')");
TestUtils.assertResultSetEqual(
statement.executeQuery("select id from pipes where creation_time > 0"),
"id,",
Expand Down Expand Up @@ -980,8 +979,7 @@ public void testDBAuth() throws SQLException {
final Statement adminStmt = adminCon.createStatement()) {
adminStmt.execute("create user test 'password123456'");
adminStmt.execute("create database db");
adminStmt.execute(
"create pipe a2b with source('double-living'='true') with sink ('sink'='write-back-sink')");
adminStmt.execute("create pipe a2b with sink ('sink'='do-nothing-sink')");
}

try (final Connection userCon =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,6 @@ public void setUp() throws Exception {
EnvFactory.getEnv()
.getConfig()
.getCommonConfig()
.setSubscriptionEnabled(true)
.setPipeMemoryManagementEnabled(false)
.setIsPipeEnableMemoryCheck(false)
// Lower the owner-lease floor so the test can use a short lease and stay fast.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,6 @@ public void setUp() throws Exception {
.setSchemaReplicationFactor(1)
.setDataReplicationFactor(2)
.setAutoCreateSchemaEnabled(true)
.setSubscriptionEnabled(true)
.setPipeMemoryManagementEnabled(false)
.setIsPipeEnableMemoryCheck(false)
.setSubscriptionOwnerLeaseDurationMsMin(1000);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,11 +32,9 @@ public abstract class AbstractSubscriptionLocalIT extends AbstractSubscriptionIT
public void setUp() throws Exception {
super.setUp();

// enable subscription
EnvFactory.getEnv()
.getConfig()
.getCommonConfig()
.setSubscriptionEnabled(true)
.setPipeMemoryManagementEnabled(false)
.setIsPipeEnableMemoryCheck(false);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,6 @@ public void setUp() throws Exception {
EnvFactory.getEnv()
.getConfig()
.getCommonConfig()
.setSubscriptionEnabled(true)
.setPipeMemoryManagementEnabled(false)
.setIsPipeEnableMemoryCheck(false)
// Lower the owner-lease floor so the drain test can use a short lease and stay fast.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,6 @@
import org.apache.iotdb.commons.pipe.agent.plugin.builtin.BuiltinPipePlugin;
import org.apache.iotdb.commons.pipe.agent.plugin.builtin.sink.donothing.DoNothingSink;
import org.apache.iotdb.commons.pipe.agent.plugin.constructor.PipeSinkConstructor;
import org.apache.iotdb.confignode.manager.pipe.sink.protocol.IoTDBConfigRegionAirGapSink;
import org.apache.iotdb.confignode.manager.pipe.sink.protocol.IoTDBConfigRegionSink;
import org.apache.iotdb.pipe.api.PipeConnector;

Expand All @@ -41,9 +40,6 @@ protected void initConstructors() {
pluginConstructors.put(
BuiltinPipePlugin.IOTDB_THRIFT_ASYNC_CONNECTOR.getPipePluginName(),
IoTDBConfigRegionSink::new);
pluginConstructors.put(
BuiltinPipePlugin.IOTDB_AIR_GAP_CONNECTOR.getPipePluginName(),
IoTDBConfigRegionAirGapSink::new);
pluginConstructors.put(
BuiltinPipePlugin.DO_NOTHING_CONNECTOR.getPipePluginName(), DoNothingSink::new);

Expand All @@ -55,8 +51,6 @@ protected void initConstructors() {
BuiltinPipePlugin.IOTDB_THRIFT_SYNC_SINK.getPipePluginName(), IoTDBConfigRegionSink::new);
pluginConstructors.put(
BuiltinPipePlugin.IOTDB_THRIFT_ASYNC_SINK.getPipePluginName(), IoTDBConfigRegionSink::new);
pluginConstructors.put(
BuiltinPipePlugin.IOTDB_AIR_GAP_SINK.getPipePluginName(), IoTDBConfigRegionAirGapSink::new);
pluginConstructors.put(
BuiltinPipePlugin.DO_NOTHING_SINK.getPipePluginName(), DoNothingSink::new);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2573,6 +2573,8 @@ private DataNodePipeMessages() {}
public static final String MESSAGE_TRANSFER_FILE_ARG_ERROR_RESULT_STATUS_ARG_E565D9FD =
"Transfer file %s error, result status %s.";

public static final String EXCEPTION_THE_PARAMETER_ARG_CANNOT_BE_SET_TO_FALSE_111C3408 =
"The parameter %s cannot be set to false.";
public static final String EXCEPTION_LEGACY_PIPE_RECEIVER_REQUIRES_A_LOGGED_IN_SESSION_D96219BF =
"Legacy pipe receiver requires a logged-in session.";
}
Original file line number Diff line number Diff line change
Expand Up @@ -2402,6 +2402,8 @@ private DataNodePipeMessages() {}
public static final String MESSAGE_TRANSFER_FILE_ARG_ERROR_RESULT_STATUS_ARG_E565D9FD =
"传输文件 %s 出错,结果状态为 %s。";

public static final String EXCEPTION_THE_PARAMETER_ARG_CANNOT_BE_SET_TO_FALSE_111C3408 =
"参数 %s 不能设置为 false。";
public static final String EXCEPTION_LEGACY_PIPE_RECEIVER_REQUIRES_A_LOGGED_IN_SESSION_D96219BF =
"Legacy pipe receiver 需要已登录的 session。";
}
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@

package org.apache.iotdb.db.conf;

import org.apache.iotdb.commons.conf.CommonConfig;
import org.apache.iotdb.commons.conf.ConfigurationFileUtils;
import org.apache.iotdb.commons.conf.TrimProperties;
import org.apache.iotdb.commons.memory.MemoryConfig;
Expand Down Expand Up @@ -521,9 +522,7 @@ private void initQueryEngineMemoryAllocate(
long maxMemoryAvailable = queryEngineMemoryManager.getTotalMemorySizeInBytes();
String queryMemoryAllocateProportion =
properties.getProperty("chunk_timeseriesmeta_free_memory_proportion");
boolean subscriptionEnabled =
Boolean.parseBoolean(
properties.getProperty("subscription_enabled", Boolean.TRUE.toString()));
boolean subscriptionEnabled = CommonConfig.SUBSCRIPTION_ENABLED;
final int[] queryMemoryProportions;
try {
queryMemoryProportions =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,18 +21,10 @@

import org.apache.iotdb.commons.pipe.agent.plugin.builtin.BuiltinPipePlugin;
import org.apache.iotdb.commons.pipe.agent.plugin.builtin.processor.donothing.DoNothingProcessor;
import org.apache.iotdb.commons.pipe.agent.plugin.builtin.processor.throwing.ThrowingExceptionProcessor;
import org.apache.iotdb.commons.pipe.agent.plugin.constructor.PipeProcessorConstructor;
import org.apache.iotdb.commons.pipe.agent.plugin.meta.DataNodePipePluginMetaKeeper;
import org.apache.iotdb.db.pipe.processor.aggregate.AggregateProcessor;
import org.apache.iotdb.db.pipe.processor.aggregate.operator.processor.StandardStatisticsOperatorProcessor;
import org.apache.iotdb.db.pipe.processor.aggregate.window.processor.TumblingWindowingProcessor;
import org.apache.iotdb.db.pipe.processor.downsampling.changing.ChangingValueSamplingProcessor;
import org.apache.iotdb.db.pipe.processor.downsampling.sdt.SwingingDoorTrendingSamplingProcessor;
import org.apache.iotdb.db.pipe.processor.downsampling.tumbling.TumblingTimeSamplingProcessor;
import org.apache.iotdb.db.pipe.processor.iotconsensusv2.IoTConsensusV2Processor;
import org.apache.iotdb.db.pipe.processor.schemachange.RenameDatabaseProcessor;
import org.apache.iotdb.db.pipe.processor.twostage.plugin.TwoStageCountProcessor;

class PipeDataRegionProcessorConstructor extends PipeProcessorConstructor {

Expand All @@ -44,28 +36,6 @@ class PipeDataRegionProcessorConstructor extends PipeProcessorConstructor {
protected void initConstructors() {
pluginConstructors.put(
BuiltinPipePlugin.DO_NOTHING_PROCESSOR.getPipePluginName(), DoNothingProcessor::new);
pluginConstructors.put(
BuiltinPipePlugin.TUMBLING_TIME_SAMPLING_PROCESSOR.getPipePluginName(),
TumblingTimeSamplingProcessor::new);
pluginConstructors.put(
BuiltinPipePlugin.SDT_SAMPLING_PROCESSOR.getPipePluginName(),
SwingingDoorTrendingSamplingProcessor::new);
pluginConstructors.put(
BuiltinPipePlugin.CHANGING_VALUE_SAMPLING_PROCESSOR.getPipePluginName(),
ChangingValueSamplingProcessor::new);
pluginConstructors.put(
BuiltinPipePlugin.THROWING_EXCEPTION_PROCESSOR.getPipePluginName(),
ThrowingExceptionProcessor::new);
pluginConstructors.put(
BuiltinPipePlugin.AGGREGATE_PROCESSOR.getPipePluginName(), AggregateProcessor::new);
pluginConstructors.put(
BuiltinPipePlugin.STANDARD_STATISTICS_PROCESSOR.getPipePluginName(),
StandardStatisticsOperatorProcessor::new);
pluginConstructors.put(
BuiltinPipePlugin.TUMBLING_WINDOWING_PROCESSOR.getPipePluginName(),
TumblingWindowingProcessor::new);
pluginConstructors.put(
BuiltinPipePlugin.COUNT_POINT_PROCESSOR.getPipePluginName(), TwoStageCountProcessor::new);
pluginConstructors.put(
BuiltinPipePlugin.IOT_CONSENSUS_V2_PROCESSOR.getPipePluginName(),
IoTConsensusV2Processor::new);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,11 +23,8 @@
import org.apache.iotdb.commons.pipe.agent.plugin.builtin.sink.donothing.DoNothingSink;
import org.apache.iotdb.commons.pipe.agent.plugin.constructor.PipeSinkConstructor;
import org.apache.iotdb.commons.pipe.agent.plugin.meta.DataNodePipePluginMetaKeeper;
import org.apache.iotdb.db.pipe.sink.protocol.airgap.IoTDBDataRegionAirGapSink;
import org.apache.iotdb.db.pipe.sink.protocol.iotconsensusv2.IoTConsensusV2AsyncSink;
import org.apache.iotdb.db.pipe.sink.protocol.legacy.IoTDBLegacyPipeSink;
import org.apache.iotdb.db.pipe.sink.protocol.opcda.OpcDaSink;
import org.apache.iotdb.db.pipe.sink.protocol.opcua.OpcUaSink;
import org.apache.iotdb.db.pipe.sink.protocol.thrift.async.IoTDBDataRegionAsyncSink;
import org.apache.iotdb.db.pipe.sink.protocol.thrift.sync.IoTDBDataRegionSyncSink;
import org.apache.iotdb.db.pipe.sink.protocol.websocket.WebSocketSink;
Expand Down Expand Up @@ -63,13 +60,8 @@ protected void initConstructors() {
pluginConstructors.put(
BuiltinPipePlugin.IOTDB_LEGACY_PIPE_CONNECTOR.getPipePluginName(),
IoTDBLegacyPipeSink::new);
pluginConstructors.put(
BuiltinPipePlugin.IOTDB_AIR_GAP_CONNECTOR.getPipePluginName(),
IoTDBDataRegionAirGapSink::new);
pluginConstructors.put(
BuiltinPipePlugin.WEBSOCKET_CONNECTOR.getPipePluginName(), WebSocketSink::new);
pluginConstructors.put(BuiltinPipePlugin.OPC_UA_CONNECTOR.getPipePluginName(), OpcUaSink::new);
pluginConstructors.put(BuiltinPipePlugin.OPC_DA_CONNECTOR.getPipePluginName(), OpcDaSink::new);
pluginConstructors.put(
BuiltinPipePlugin.DO_NOTHING_CONNECTOR.getPipePluginName(), DoNothingSink::new);
pluginConstructors.put(
Expand All @@ -86,12 +78,8 @@ protected void initConstructors() {
IoTDBDataRegionAsyncSink::new);
pluginConstructors.put(
BuiltinPipePlugin.IOTDB_LEGACY_PIPE_SINK.getPipePluginName(), IoTDBLegacyPipeSink::new);
pluginConstructors.put(
BuiltinPipePlugin.IOTDB_AIR_GAP_SINK.getPipePluginName(), IoTDBDataRegionAirGapSink::new);
pluginConstructors.put(
BuiltinPipePlugin.WEBSOCKET_SINK.getPipePluginName(), WebSocketSink::new);
pluginConstructors.put(BuiltinPipePlugin.OPC_UA_SINK.getPipePluginName(), OpcUaSink::new);
pluginConstructors.put(BuiltinPipePlugin.OPC_DA_SINK.getPipePluginName(), OpcDaSink::new);
pluginConstructors.put(
BuiltinPipePlugin.DO_NOTHING_SINK.getPipePluginName(), DoNothingSink::new);
pluginConstructors.put(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,6 @@
import org.apache.iotdb.commons.pipe.agent.plugin.builtin.BuiltinPipePlugin;
import org.apache.iotdb.commons.pipe.agent.plugin.builtin.sink.donothing.DoNothingSink;
import org.apache.iotdb.commons.pipe.agent.plugin.constructor.PipeSinkConstructor;
import org.apache.iotdb.db.pipe.sink.protocol.airgap.IoTDBSchemaRegionAirGapSink;
import org.apache.iotdb.db.pipe.sink.protocol.thrift.sync.IoTDBSchemaRegionSink;
import org.apache.iotdb.pipe.api.PipeConnector;

Expand All @@ -41,9 +40,6 @@ protected void initConstructors() {
pluginConstructors.put(
BuiltinPipePlugin.IOTDB_THRIFT_ASYNC_CONNECTOR.getPipePluginName(),
IoTDBSchemaRegionSink::new);
pluginConstructors.put(
BuiltinPipePlugin.IOTDB_AIR_GAP_CONNECTOR.getPipePluginName(),
IoTDBSchemaRegionAirGapSink::new);
pluginConstructors.put(
BuiltinPipePlugin.DO_NOTHING_CONNECTOR.getPipePluginName(), DoNothingSink::new);

Expand All @@ -55,8 +51,6 @@ protected void initConstructors() {
BuiltinPipePlugin.IOTDB_THRIFT_SYNC_SINK.getPipePluginName(), IoTDBSchemaRegionSink::new);
pluginConstructors.put(
BuiltinPipePlugin.IOTDB_THRIFT_ASYNC_SINK.getPipePluginName(), IoTDBSchemaRegionSink::new);
pluginConstructors.put(
BuiltinPipePlugin.IOTDB_AIR_GAP_SINK.getPipePluginName(), IoTDBSchemaRegionAirGapSink::new);
pluginConstructors.put(
BuiltinPipePlugin.DO_NOTHING_SINK.getPipePluginName(), DoNothingSink::new);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -78,7 +78,7 @@ public synchronized String register(
connectorKey);

final boolean isDataRegionSink = isDataRegionSink(environment.getRegionId());
final int sinkNum = calculateSinkSubtaskNum(pipeSinkParameters, isDataRegionSink, connectorKey);
final int sinkNum = calculateSinkSubtaskNum(pipeSinkParameters, isDataRegionSink);
boolean realTimeFirst = false;
final String attributeSortedString =
generateAttributeSortedString(pipeSinkParameters, environment.getRegionId());
Expand Down Expand Up @@ -260,8 +260,7 @@ public synchronized boolean hasRegisteredSubtasks(

public static int calculateSinkSubtaskNum(
final PipeParameters pipeSinkParameters, final int regionId) {
final String connectorKey = getConnectorKey(pipeSinkParameters);
return calculateSinkSubtaskNum(pipeSinkParameters, isDataRegionSink(regionId), connectorKey);
return calculateSinkSubtaskNum(pipeSinkParameters, isDataRegionSink(regionId));
}

public static String generateAttributeSortedString(
Expand All @@ -285,9 +284,7 @@ private static boolean isDataRegionSink(final int regionId) {
}

private static int calculateSinkSubtaskNum(
final PipeParameters pipeSinkParameters,
final boolean isDataRegionSink,
final String connectorKey) {
final PipeParameters pipeSinkParameters, final boolean isDataRegionSink) {
if (!isDataRegionSink) {
// Do not allow parallel tasks for schema region connectors to avoid the potential disorder of
// the schema region data transfer.
Expand All @@ -300,9 +297,7 @@ private static int calculateSinkSubtaskNum(
Arrays.asList(
PipeSinkConstant.CONNECTOR_IOTDB_PARALLEL_TASKS_KEY,
PipeSinkConstant.SINK_IOTDB_PARALLEL_TASKS_KEY),
PipeSinkConstant.SINGLE_THREAD_DEFAULT_SINK.contains(connectorKey)
? 1
: PipeSinkConstant.CONNECTOR_IOTDB_PARALLEL_TASKS_DEFAULT_VALUE);
PipeSinkConstant.CONNECTOR_IOTDB_PARALLEL_TASKS_DEFAULT_VALUE);
}

private static String generateAttributeSortedString(
Expand Down
Loading
Loading