Skip to content

Commit 312e047

Browse files
committed
Optimized the hint for subscription (#17115)
* ams * y * fix
1 parent 230e3a8 commit 312e047

14 files changed

Lines changed: 94 additions & 99 deletions

File tree

iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/agent/plugin/PipeConfigNodePluginAgent.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,7 @@ public PipeConfigNodePluginAgent(PipePluginMetaKeeper pipePluginMetaKeeper) {
3232
}
3333

3434
@Override
35-
protected PipeSourceConstructor createPipeExtractorConstructor(
35+
protected PipeSourceConstructor createPipeSourceConstructor(
3636
PipePluginMetaKeeper pipePluginMetaKeeper) {
3737
return new PipeConfigRegionSourceConstructor();
3838
}
@@ -44,7 +44,7 @@ protected PipeProcessorConstructor createPipeProcessorConstructor(
4444
}
4545

4646
@Override
47-
protected PipeSinkConstructor createPipeConnectorConstructor(
47+
protected PipeSinkConstructor createPipeSinkConstructor(
4848
PipePluginMetaKeeper pipePluginMetaKeeper) {
4949
return new PipeConfigRegionSinkConstructor();
5050
}

iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/sink/protocol/IoTDBConfigRegionSink.java

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -183,9 +183,7 @@ private void doTransfer(final PipeConfigRegionWritePlanEvent pipeConfigRegionWri
183183
true);
184184
}
185185

186-
if (LOGGER.isDebugEnabled()) {
187-
LOGGER.debug("Successfully transferred config event {}.", pipeConfigRegionWritePlanEvent);
188-
}
186+
LOGGER.info("Successfully transferred config event {}.", pipeConfigRegionWritePlanEvent);
189187
}
190188

191189
private void doTransferWrapper(final PipeConfigRegionSnapshotEvent pipeConfigRegionSnapshotEvent)

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/plugin/PipeDataNodePluginAgent.java

Lines changed: 6 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -212,19 +212,18 @@ public void validate(
212212
pipeName, extractorAttributes, processorAttributes, connectorAttributes);
213213
}
214214

215-
public void validateExtractor(Map<String, String> extractorAttributes) throws Exception {
216-
dataRegionAgent.validateExtractor(extractorAttributes);
217-
schemaRegionAgent.validateExtractor(extractorAttributes);
215+
public void validateSource(Map<String, String> sourceAttributes) throws Exception {
216+
dataRegionAgent.validateSource(sourceAttributes);
217+
schemaRegionAgent.validateSource(sourceAttributes);
218218
}
219219

220220
public void validateProcessor(Map<String, String> processorAttributes) throws Exception {
221221
dataRegionAgent.validateProcessor(processorAttributes);
222222
schemaRegionAgent.validateProcessor(processorAttributes);
223223
}
224224

225-
public void validateConnector(String pipeName, Map<String, String> connectorAttributes)
226-
throws Exception {
227-
dataRegionAgent.validateConnector(pipeName, connectorAttributes);
228-
schemaRegionAgent.validateConnector(pipeName, connectorAttributes);
225+
public void validateSource(String pipeName, Map<String, String> sinkAttributes) throws Exception {
226+
dataRegionAgent.validateSink(pipeName, sinkAttributes);
227+
schemaRegionAgent.validateSink(pipeName, sinkAttributes);
229228
}
230229
}

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/plugin/dataregion/PipeDataRegionPluginAgent.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,7 @@ public PipeDataRegionPluginAgent(DataNodePipePluginMetaKeeper pipePluginMetaKeep
3333
}
3434

3535
@Override
36-
protected PipeSourceConstructor createPipeExtractorConstructor(
36+
protected PipeSourceConstructor createPipeSourceConstructor(
3737
PipePluginMetaKeeper pipePluginMetaKeeper) {
3838
return new PipeDataRegionSourceConstructor((DataNodePipePluginMetaKeeper) pipePluginMetaKeeper);
3939
}
@@ -46,7 +46,7 @@ protected PipeProcessorConstructor createPipeProcessorConstructor(
4646
}
4747

4848
@Override
49-
protected PipeSinkConstructor createPipeConnectorConstructor(
49+
protected PipeSinkConstructor createPipeSinkConstructor(
5050
PipePluginMetaKeeper pipePluginMetaKeeper) {
5151
return new PipeDataRegionSinkConstructor((DataNodePipePluginMetaKeeper) pipePluginMetaKeeper);
5252
}

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/plugin/schemaregion/PipeSchemaRegionPluginAgent.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,7 @@ public PipeSchemaRegionPluginAgent(PipePluginMetaKeeper pipePluginMetaKeeper) {
3232
}
3333

3434
@Override
35-
protected PipeSourceConstructor createPipeExtractorConstructor(
35+
protected PipeSourceConstructor createPipeSourceConstructor(
3636
PipePluginMetaKeeper pipePluginMetaKeeper) {
3737
return new PipeSchemaRegionSourceConstructor();
3838
}
@@ -44,7 +44,7 @@ protected PipeProcessorConstructor createPipeProcessorConstructor(
4444
}
4545

4646
@Override
47-
protected PipeSinkConstructor createPipeConnectorConstructor(
47+
protected PipeSinkConstructor createPipeSinkConstructor(
4848
PipePluginMetaKeeper pipePluginMetaKeeper) {
4949
return new PipeSchemaRegionSinkConstructor();
5050
}

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/builder/PipeDataNodeTaskBuilder.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -124,7 +124,7 @@ public PipeDataNodeTask build() {
124124
blendUserAndSystemParameters(pipeStaticMeta.getProcessorParameters()),
125125
regionId,
126126
sourceStage.getEventSupplier(),
127-
sinkStage.getPipeConnectorPendingQueue(),
127+
sinkStage.getPipeSinkPendingQueue(),
128128
PROCESSOR_EXECUTOR,
129129
pipeTaskMeta,
130130
pipeStaticMeta

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/stage/PipeTaskSinkStage.java

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,7 @@ public class PipeTaskSinkStage extends PipeTaskStage {
3434

3535
protected final String pipeName;
3636
protected final long creationTime;
37-
protected final PipeParameters pipeConnectorParameters;
37+
protected final PipeParameters pipeSinkParameters;
3838
protected final int regionId;
3939
protected final Supplier<? extends PipeSinkSubtaskExecutor> executor;
4040

@@ -43,12 +43,12 @@ public class PipeTaskSinkStage extends PipeTaskStage {
4343
public PipeTaskSinkStage(
4444
String pipeName,
4545
long creationTime,
46-
PipeParameters pipeConnectorParameters,
46+
PipeParameters pipeSinkParameters,
4747
int regionId,
4848
Supplier<? extends PipeSinkSubtaskExecutor> executor) {
4949
this.pipeName = pipeName;
5050
this.creationTime = creationTime;
51-
this.pipeConnectorParameters = pipeConnectorParameters;
51+
this.pipeSinkParameters = pipeSinkParameters;
5252
this.regionId = regionId;
5353
this.executor = executor;
5454

@@ -60,7 +60,7 @@ protected void registerSubtask() {
6060
PipeSinkSubtaskManager.instance()
6161
.register(
6262
executor,
63-
pipeConnectorParameters,
63+
pipeSinkParameters,
6464
new PipeTaskSinkRuntimeEnvironment(pipeName, creationTime, regionId));
6565
}
6666

@@ -85,7 +85,7 @@ public void dropSubtask() throws PipeException {
8585
.deregister(pipeName, creationTime, regionId, connectorSubtaskId);
8686
}
8787

88-
public UnboundedBlockingPendingQueue<Event> getPipeConnectorPendingQueue() {
89-
return PipeSinkSubtaskManager.instance().getPipeConnectorPendingQueue(connectorSubtaskId);
88+
public UnboundedBlockingPendingQueue<Event> getPipeSinkPendingQueue() {
89+
return PipeSinkSubtaskManager.instance().getPipeSinkPendingQueue(connectorSubtaskId);
9090
}
9191
}

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtaskLifeCycle.java

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -39,9 +39,9 @@ public class PipeSinkSubtaskLifeCycle implements AutoCloseable {
3939
protected int registeredTaskCount;
4040

4141
public PipeSinkSubtaskLifeCycle(
42-
PipeSinkSubtaskExecutor executor,
43-
PipeSinkSubtask subtask,
44-
UnboundedBlockingPendingQueue<Event> pendingQueue) {
42+
final PipeSinkSubtaskExecutor executor,
43+
final PipeSinkSubtask subtask,
44+
final UnboundedBlockingPendingQueue<Event> pendingQueue) {
4545
this.executor = executor;
4646
this.subtask = subtask;
4747
this.pendingQueue = pendingQueue;

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtaskManager.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -240,7 +240,7 @@ public synchronized void stop(final String attributeSortedString) {
240240
}
241241
}
242242

243-
public UnboundedBlockingPendingQueue<Event> getPipeConnectorPendingQueue(
243+
public UnboundedBlockingPendingQueue<Event> getPipeSinkPendingQueue(
244244
final String attributeSortedString) {
245245
if (!attributeSortedString2SubtaskLifeCycleMap.containsKey(attributeSortedString)) {
246246
throw new PipeException(

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeRemainingEventAndTimeOperator.java

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,7 @@
4343
public class PipeDataNodeRemainingEventAndTimeOperator extends PipeRemainingOperator {
4444

4545
// Calculate from schema region extractors directly for it requires less computation
46-
private final Set<IoTDBSchemaRegionSource> schemaRegionExtractors =
46+
private final Set<IoTDBSchemaRegionSource> schemaRegionSources =
4747
Collections.newSetFromMap(new ConcurrentHashMap<>());
4848

4949
private final AtomicInteger insertNodeEventCount = new AtomicInteger(0);
@@ -105,7 +105,7 @@ public long getRemainingNonHeartbeatEvents() {
105105
tsfileEventCount.get()
106106
+ rawTabletEventCount.get()
107107
+ insertNodeEventCount.get()
108-
+ schemaRegionExtractors.stream()
108+
+ schemaRegionSources.stream()
109109
.map(IoTDBSchemaRegionSource::getUnTransferredEventCount)
110110
.reduce(Long::sum)
111111
.orElse(0L);
@@ -157,7 +157,7 @@ public double getRemainingTime() {
157157
}
158158

159159
final long totalSchemaRegionWriteEventCount =
160-
schemaRegionExtractors.stream()
160+
schemaRegionSources.stream()
161161
.map(IoTDBSchemaRegionSource::getUnTransferredEventCount)
162162
.reduce(Long::sum)
163163
.orElse(0L);
@@ -192,8 +192,8 @@ public double getRemainingTime() {
192192

193193
//////////////////////////// Register & deregister (pipe integration) ////////////////////////////
194194

195-
void register(final IoTDBSchemaRegionSource extractor) {
196-
schemaRegionExtractors.add(extractor);
195+
void register(final IoTDBSchemaRegionSource source) {
196+
schemaRegionSources.add(source);
197197
}
198198

199199
//////////////////////////// Rate ////////////////////////////

0 commit comments

Comments
 (0)