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
Original file line number Diff line number Diff line change
Expand Up @@ -204,14 +204,6 @@ public final class ManagerMessages {
"Failed to sync template {} extension info to DataNode {}";
public static final String FAILED_TO_SYNC_TOPIC_META_RESULT_STATUS =
"Failed to sync topic meta. Result status: {}.";
public static final String FAILED_TO_UNBIND_FROM_PIPE_CONFIG_REGION_CONNECTOR_METRICS_CONNECTOR =
"Failed to unbind from pipe config region connector metrics, connector map not empty";
public static final String FAILED_TO_UNBIND_FROM_PIPE_CONFIG_REGION_EXTRACTOR_METRICS_EXTRACTOR =
"Failed to unbind from pipe config region extractor metrics, extractor map not empty";
public static final String FAILED_TO_UNBIND_FROM_PIPE_REMAINING_TIME_METRICS_REMAININGTIMEOPERATOR_MAP =
"Failed to unbind from pipe remaining time metrics, RemainingTimeOperator map not empty";
public static final String FAILED_TO_UNBIND_FROM_PIPE_TEMPORARY_META_METRICS_PIPETEMPORARYMETA_MAP =
"Failed to unbind from pipe temporary meta metrics, PipeTemporaryMeta map not empty";
public static final String FAILED_TO_UPDATE_PIPE_PROCEDURE_TIMER_PIPEPROCEDURE_DOES_NOT_EXIST =
"Failed to update pipe procedure timer, PipeProcedure({}) does not exist";
public static final String FAILED_TO_UPDATE_THE_LAST_EXECUTION_TIME_OF_CQ_BECAUSE =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -202,14 +202,6 @@ public final class ManagerMessages {
"将模板 {} 的扩展信息同步到 DataNode {} 失败";
public static final String FAILED_TO_SYNC_TOPIC_META_RESULT_STATUS =
"同步 topic 元数据失败。结果状态:{}。";
public static final String FAILED_TO_UNBIND_FROM_PIPE_CONFIG_REGION_CONNECTOR_METRICS_CONNECTOR =
"从 pipe config region connector 指标解绑失败,connector map 不为空";
public static final String FAILED_TO_UNBIND_FROM_PIPE_CONFIG_REGION_EXTRACTOR_METRICS_EXTRACTOR =
"从 pipe config region extractor 指标解绑失败,extractor map 不为空";
public static final String FAILED_TO_UNBIND_FROM_PIPE_REMAINING_TIME_METRICS_REMAININGTIMEOPERATOR_MAP =
"从 pipe remaining time 指标解绑失败,RemainingTimeOperator map 不为空";
public static final String FAILED_TO_UNBIND_FROM_PIPE_TEMPORARY_META_METRICS_PIPETEMPORARYMETA_MAP =
"从 pipe temporary meta 指标解绑失败,PipeTemporaryMeta map 不为空";
public static final String FAILED_TO_UPDATE_PIPE_PROCEDURE_TIMER_PIPEPROCEDURE_DOES_NOT_EXIST =
"更新 pipe procedure timer 失败,PipeProcedure({}) 不存在";
public static final String FAILED_TO_UPDATE_THE_LAST_EXECUTION_TIME_OF_CQ_BECAUSE =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@ public class PipeConfigNodeRemainingTimeMetrics implements IMetricSet {
//////////////////////////// bindTo & unbindFrom (metric framework) ////////////////////////////

@Override
public void bindTo(final AbstractMetricService metricService) {
public synchronized void bindTo(final AbstractMetricService metricService) {
this.metricService = metricService;
ImmutableSet.copyOf(remainingTimeOperatorMap.keySet()).forEach(this::createMetrics);
}
Expand All @@ -74,13 +74,11 @@ private void createAutoGauge(final String pipeID) {
}

@Override
public void unbindFrom(final AbstractMetricService metricService) {
ImmutableSet.copyOf(remainingTimeOperatorMap.keySet()).forEach(this::deregister);
if (!remainingTimeOperatorMap.isEmpty()) {
LOGGER.warn(
ManagerMessages
.FAILED_TO_UNBIND_FROM_PIPE_REMAINING_TIME_METRICS_REMAININGTIMEOPERATOR_MAP);
}
public synchronized void unbindFrom(final AbstractMetricService metricService) {
// Keep the operators registered: they hold the states of the pipes and register only once, so
// a metric service restart must be able to bind them again.
// Synchronized with the (de)registrations, which may remove a registration being unbound.
ImmutableSet.copyOf(remainingTimeOperatorMap.keySet()).forEach(this::removeMetrics);
}

private void removeMetrics(final String pipeID) {
Expand All @@ -96,12 +94,11 @@ private void removeAutoGauge(final String pipeID) {
operator.getPipeName(),
Tag.CREATION_TIME.toString(),
String.valueOf(operator.getCreationTime()));
remainingTimeOperatorMap.remove(pipeID);
}

//////////////////////////// register & deregister (pipe integration) ////////////////////////////

public void register(final IoTDBConfigRegionSource extractor) {
public synchronized void register(final IoTDBConfigRegionSource extractor) {
// The metric is global thus the regionId is omitted
final String pipeID = extractor.getPipeName() + "_" + extractor.getCreationTime();
remainingTimeOperatorMap
Expand Down Expand Up @@ -132,7 +129,7 @@ public void freezeRate(final String pipeID) {
remainingTimeOperatorMap.get(pipeID).freezeRate(true);
}

public void deregister(final String pipeID) {
public synchronized void deregister(final String pipeID) {
if (!remainingTimeOperatorMap.containsKey(pipeID)) {
LOGGER.warn(
ManagerMessages
Expand All @@ -143,6 +140,7 @@ public void deregister(final String pipeID) {
if (Objects.nonNull(metricService)) {
removeMetrics(pipeID);
}
remainingTimeOperatorMap.remove(pipeID);
}

public void markRegionCommit(final String pipeID, final boolean isDataRegion) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,7 @@ public class PipeTemporaryMetaInCoordinatorMetrics implements IMetricSet {
//////////////////////////// bindTo & unbindFrom (metric framework) ////////////////////////////

@Override
public void bindTo(final AbstractMetricService metricService) {
public synchronized void bindTo(final AbstractMetricService metricService) {
this.metricService = metricService;
ImmutableSet.copyOf(pipeTemporaryMetaMap.keySet()).forEach(this::createMetrics);
}
Expand Down Expand Up @@ -91,12 +91,10 @@ private void createAutoGauge(final String pipeID) {
}

@Override
public void unbindFrom(final AbstractMetricService metricService) {
ImmutableSet.copyOf(pipeTemporaryMetaMap.keySet()).forEach(this::deregister);
if (!pipeTemporaryMetaMap.isEmpty()) {
LOGGER.warn(
ManagerMessages.FAILED_TO_UNBIND_FROM_PIPE_TEMPORARY_META_METRICS_PIPETEMPORARYMETA_MAP);
}
public synchronized void unbindFrom(final AbstractMetricService metricService) {
// Keep the pipes registered so that a metric service restart can bind them again.
// Synchronized with the (de)registrations, which may remove a registration being unbound.
ImmutableSet.copyOf(pipeTemporaryMetaMap.keySet()).forEach(this::removeMetrics);
}

private void removeMetrics(final String pipeID) {
Expand All @@ -119,12 +117,11 @@ private void removeAutoGauge(final String pipeID) {
pipeNameAndCreationTime[0],
Tag.CREATION_TIME.toString(),
pipeNameAndCreationTime[1]);
pipeTemporaryMetaMap.remove(pipeID);
}

//////////////////////////// register & deregister (pipe integration) ////////////////////////////

public void register(final PipeMeta pipeMeta) {
public synchronized void register(final PipeMeta pipeMeta) {
final String taskID =
pipeMeta.getStaticMeta().getPipeName() + "_" + pipeMeta.getStaticMeta().getCreationTime();
pipeTemporaryMetaMap.putIfAbsent(
Expand All @@ -134,7 +131,7 @@ public void register(final PipeMeta pipeMeta) {
}
}

public void deregister(final String pipeID) {
public synchronized void deregister(final String pipeID) {
if (!pipeTemporaryMetaMap.containsKey(pipeID)) {
LOGGER.warn(
ManagerMessages
Expand All @@ -145,6 +142,7 @@ public void deregister(final String pipeID) {
if (Objects.nonNull(metricService)) {
removeMetrics(pipeID);
}
pipeTemporaryMetaMap.remove(pipeID);
}

public void handleTemporaryMetaChanges(final Iterable<PipeMeta> pipeMetaList) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@ public class PipeConfigRegionSinkMetrics implements IMetricSet {
//////////////////////////// bindTo & unbindFrom (metric framework) ////////////////////////////

@Override
public void bindTo(final AbstractMetricService metricService) {
public synchronized void bindTo(final AbstractMetricService metricService) {
this.metricService = metricService;
ImmutableSet.copyOf(subtaskMap.keySet()).forEach(this::createMetrics);
}
Expand All @@ -74,12 +74,11 @@ private void createRate(final String taskID) {
}

@Override
public void unbindFrom(final AbstractMetricService metricService) {
ImmutableSet.copyOf(subtaskMap.keySet()).forEach(this::deregister);
if (!subtaskMap.isEmpty()) {
LOGGER.warn(
ManagerMessages.FAILED_TO_UNBIND_FROM_PIPE_CONFIG_REGION_CONNECTOR_METRICS_CONNECTOR);
}
public synchronized void unbindFrom(final AbstractMetricService metricService) {
// Keep the subtasks registered: they register only once, so a metric service restart
// must be able to bind them again.
// Synchronized with the (de)registrations, which may remove a registration being unbound.
ImmutableSet.copyOf(subtaskMap.keySet()).forEach(this::removeMetrics);
}

private void removeMetrics(final String taskID) {
Expand All @@ -101,15 +100,15 @@ private void removeRate(final String taskID) {

//////////////////////////// register & deregister (pipe integration) ////////////////////////////

public void register(final PipeConfigNodeSubtask pipeConfigNodeSubtask) {
public synchronized void register(final PipeConfigNodeSubtask pipeConfigNodeSubtask) {
final String taskID = pipeConfigNodeSubtask.getTaskID();
subtaskMap.putIfAbsent(taskID, pipeConfigNodeSubtask);
if (Objects.nonNull(metricService)) {
createMetrics(taskID);
}
}

public void deregister(final String taskID) {
public synchronized void deregister(final String taskID) {
if (!subtaskMap.containsKey(taskID)) {
LOGGER.warn(ManagerMessages.FAILED_TO_DEREGISTER_PIPE_CONFIG_REGION_CONNECTOR, taskID);
return;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@ public class PipeConfigRegionSourceMetrics implements IMetricSet {
//////////////////////////// bindTo & unbindFrom (metric framework) ////////////////////////////

@Override
public void bindTo(final AbstractMetricService metricService) {
public synchronized void bindTo(final AbstractMetricService metricService) {
this.metricService = metricService;
ImmutableSet.copyOf(extractorMap.keySet()).forEach(this::createMetrics);
}
Expand All @@ -70,12 +70,11 @@ private void createAutoGauge(final String taskID) {
}

@Override
public void unbindFrom(final AbstractMetricService metricService) {
ImmutableSet.copyOf(extractorMap.keySet()).forEach(this::deregister);
if (!extractorMap.isEmpty()) {
LOGGER.warn(
ManagerMessages.FAILED_TO_UNBIND_FROM_PIPE_CONFIG_REGION_EXTRACTOR_METRICS_EXTRACTOR);
}
public synchronized void unbindFrom(final AbstractMetricService metricService) {
// Keep the extractors registered: they register only once, so a metric service restart
// must be able to bind them again.
// Synchronized with the (de)registrations, which may remove a registration being unbound.
ImmutableSet.copyOf(extractorMap.keySet()).forEach(this::removeMetrics);
}

private void removeMetrics(final String taskID) {
Expand All @@ -96,15 +95,15 @@ private void removeAutoGauge(final String taskID) {

//////////////////////////// pipe integration ////////////////////////////

public void register(final IoTDBConfigRegionSource extractor) {
public synchronized void register(final IoTDBConfigRegionSource extractor) {
final String taskID = extractor.getTaskID();
extractorMap.putIfAbsent(taskID, extractor);
if (Objects.nonNull(metricService)) {
createMetrics(taskID);
}
}

public void deregister(final String taskID) {
public synchronized void deregister(final String taskID) {
if (!extractorMap.containsKey(taskID)) {
LOGGER.warn(ManagerMessages.FAILED_TO_DEREGISTER_PIPE_CONFIG_REGION_EXTRACTOR, taskID);
return;
Expand Down
Loading
Loading