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 @@ -145,7 +145,7 @@ public void testAutoDropInHistoricalTransferWithTimeRange() throws Exception {
TestUtils.assertDataEventuallyOnEnv(
senderEnv,
"show pipes",
"ID,CreationTime,State,PipeSource,PipeProcessor,PipeSink,ExceptionMessage,RemainingEventCount,EstimatedRemainingSeconds,",
"ID,CreationTime,State,PipeSource,PipeProcessor,PipeSink,ExceptionMessage,RemainingEventCount,EstimatedRemainingSeconds,RecentFailures,",
Collections.emptySet());
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -147,7 +147,8 @@ public void onComplete(TDataNodeHeartbeatResp heartbeatResp) {
heartbeatResp.getPipeMetaList(),
heartbeatResp.getPipeCompletedList(),
heartbeatResp.getPipeRemainingEventCountList(),
heartbeatResp.getPipeRemainingTimeList());
heartbeatResp.getPipeRemainingTimeList(),
heartbeatResp.getPipeRecentFailureList());
}
if (heartbeatResp.isSetConfirmedConfigNodeEndPoints()) {
loadManager
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -204,6 +204,7 @@ public TShowPipeResp convertToTShowPipeResp() {
canCalculateOnLocal ? -1 : temporaryMeta.getGlobalRemainingEvents());
showPipeInfo.setEstimatedRemainingTime(
canCalculateOnLocal ? -1 : temporaryMeta.getGlobalRemainingTime());
showPipeInfo.setRecentFailures(temporaryMeta.getGlobalRecentFailures());
showPipeInfoList.add(showPipeInfo);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2576,7 +2576,8 @@ public TSStatus pushHeartbeat(final int dataNodeId, final TPipeHeartbeatResp res
resp.getPipeMetaList(),
resp.getPipeCompletedList(),
resp.getPipeRemainingEventCountList(),
resp.getPipeRemainingTimeList());
resp.getPipeRemainingTimeList(),
resp.getPipeRecentFailureList());
return StatusUtils.OK;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
import org.apache.iotdb.commons.pipe.config.plugin.env.PipeTaskSourceRuntimeEnvironment;
import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
import org.apache.iotdb.commons.pipe.event.ProgressReportEvent;
import org.apache.iotdb.commons.pipe.resource.PipeResourceFailureType;
import org.apache.iotdb.confignode.manager.pipe.agent.PipeConfigNodeAgent;
import org.apache.iotdb.confignode.manager.pipe.metric.sink.PipeConfigRegionSinkMetrics;
import org.apache.iotdb.confignode.manager.pipe.source.IoTDBConfigRegionSource;
Expand Down Expand Up @@ -240,6 +241,13 @@ protected void report(final EnrichedEvent event, final PipeRuntimeException exce
PipeConfigNodeAgent.runtime().report(event, exception);
}

@Override
protected void reportResourceFailure(
final EnrichedEvent event, final PipeResourceFailureType failureType) {
PipeConfigNodeAgent.task()
.recordPipeResourceFailure(event.getPipeName(), event.getCreationTime(), failureType);
}

//////////////////////////// APIs provided for metric framework ////////////////////////////

public String getPipeName() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeMeta;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStaticMeta;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTemporaryMetaInAgent;
import org.apache.iotdb.commons.pipe.config.PipeConfig;
import org.apache.iotdb.confignode.conf.ConfigNodeDescriptor;
import org.apache.iotdb.confignode.manager.pipe.agent.PipeConfigNodeAgent;
Expand Down Expand Up @@ -218,6 +219,7 @@ protected void collectPipeMetaListInternal(
final List<ByteBuffer> pipeMetaBinaryList = new ArrayList<>();
final List<Long> pipeRemainingEventCountList = new ArrayList<>();
final List<Double> pipeRemainingTimeList = new ArrayList<>();
final List<Map<String, Long>> pipeRecentFailureList = new ArrayList<>();
try {
for (final PipeMeta pipeMeta : pipeMetaKeeper.getPipeMetaList()) {
pipeMetaBinaryList.add(pipeMeta.serialize());
Expand All @@ -232,6 +234,8 @@ protected void collectPipeMetaListInternal(

pipeRemainingEventCountList.add(remainingEventCount);
pipeRemainingTimeList.add(estimatedRemainingTime);
pipeRecentFailureList.add(
((PipeTemporaryMetaInAgent) pipeMeta.getTemporaryMeta()).getRecentFailures());

logger.ifPresent(
l ->
Expand All @@ -248,6 +252,7 @@ protected void collectPipeMetaListInternal(
resp.setPipeMetaList(pipeMetaBinaryList);
resp.setPipeRemainingEventCountList(pipeRemainingEventCountList);
resp.setPipeRemainingTimeList(pipeRemainingTimeList);
resp.setPipeRecentFailureList(pipeRecentFailureList);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@

import java.nio.ByteBuffer;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.atomic.AtomicReference;

Expand Down Expand Up @@ -107,13 +108,15 @@ public void parseHeartbeat(
final List<ByteBuffer> pipeMetaByteBufferListFromDataNode,
/* @Nullable */ final List<Boolean> pipeCompletedListFromAgent,
/* @Nullable */ final List<Long> pipeRemainingEventCountListFromAgent,
/* @Nullable */ final List<Double> pipeRemainingTimeListFromAgent) {
/* @Nullable */ final List<Double> pipeRemainingTimeListFromAgent,
/* @Nullable */ final List<Map<String, Long>> pipeRecentFailureListFromAgent) {
pipeHeartbeatScheduler.parseHeartbeat(
dataNodeId,
new PipeHeartbeat(
pipeMetaByteBufferListFromDataNode,
pipeCompletedListFromAgent,
pipeRemainingEventCountListFromAgent,
pipeRemainingTimeListFromAgent));
pipeRemainingTimeListFromAgent,
pipeRecentFailureListFromAgent));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStaticMeta;

import java.nio.ByteBuffer;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
Expand All @@ -33,12 +34,27 @@ public class PipeHeartbeat {
private final Map<PipeStaticMeta, Boolean> isCompletedMap = new HashMap<>();
private final Map<PipeStaticMeta, Long> remainingEventCountMap = new HashMap<>();
private final Map<PipeStaticMeta, Double> remainingTimeMap = new HashMap<>();
private final Map<PipeStaticMeta, Map<String, Long>> recentFailuresMap = new HashMap<>();

public PipeHeartbeat(
final List<ByteBuffer> pipeMetaByteBufferListFromAgent,
/* @Nullable */ final List<Boolean> pipeCompletedListFromAgent,
/* @Nullable */ final List<Long> pipeRemainingEventCountListFromAgent,
/* @Nullable */ final List<Double> pipeRemainingTimeListFromAgent) {
this(
pipeMetaByteBufferListFromAgent,
pipeCompletedListFromAgent,
pipeRemainingEventCountListFromAgent,
pipeRemainingTimeListFromAgent,
null);
}

public PipeHeartbeat(
final List<ByteBuffer> pipeMetaByteBufferListFromAgent,
/* @Nullable */ final List<Boolean> pipeCompletedListFromAgent,
/* @Nullable */ final List<Long> pipeRemainingEventCountListFromAgent,
/* @Nullable */ final List<Double> pipeRemainingTimeListFromAgent,
/* @Nullable */ final List<Map<String, Long>> pipeRecentFailureListFromAgent) {
// Shall not reach here, just in case
if (Objects.isNull(pipeMetaByteBufferListFromAgent)) {
return;
Expand All @@ -63,6 +79,13 @@ public PipeHeartbeat(
Objects.nonNull(pipeRemainingTimeListFromAgent)
? pipeRemainingTimeListFromAgent.get(i)
: 0d);
recentFailuresMap.put(
pipeMeta.getStaticMeta(),
Objects.nonNull(pipeRecentFailureListFromAgent)
&& i < pipeRecentFailureListFromAgent.size()
&& Objects.nonNull(pipeRecentFailureListFromAgent.get(i))
? new HashMap<>(pipeRecentFailureListFromAgent.get(i))
: Collections.emptyMap());
}
}

Expand All @@ -86,6 +109,10 @@ public Double getRemainingTime(final PipeStaticMeta pipeStaticMeta) {
return remainingTimeMap.get(pipeStaticMeta);
}

public Map<String, Long> getRecentFailures(final PipeStaticMeta pipeStaticMeta) {
return recentFailuresMap.get(pipeStaticMeta);
}

public boolean isEmpty() {
return pipeMetaMap.isEmpty();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,7 @@ private void parseHeartbeatAndSaveMetaChangeLocally(
// Record statistics
temporaryMeta.setRemainingEvent(nodeId, pipeHeartbeat.getRemainingEventCount(staticMeta));
temporaryMeta.setRemainingTime(nodeId, pipeHeartbeat.getRemainingTime(staticMeta));
temporaryMeta.setRecentFailures(nodeId, pipeHeartbeat.getRecentFailures(staticMeta));

final Map<Integer, PipeTaskMeta> pipeTaskMetaMapFromCoordinator =
pipeMetaFromCoordinator.getRuntimeMeta().getConsensusGroupId2TaskMetaMap();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -112,7 +112,8 @@ private synchronized void heartbeat() {
resp.getPipeMetaList(),
resp.getPipeCompletedList(),
resp.getPipeRemainingEventCountList(),
resp.getPipeRemainingTimeList())));
resp.getPipeRemainingTimeList(),
resp.getPipeRecentFailureList())));

// config node heartbeat
try {
Expand All @@ -124,7 +125,8 @@ private synchronized void heartbeat() {
configNodeResp.getPipeMetaList(),
null,
configNodeResp.getPipeRemainingEventCountList(),
configNodeResp.getPipeRemainingTimeList()));
configNodeResp.getPipeRemainingTimeList(),
configNodeResp.getPipeRecentFailureList()));
} catch (final Exception e) {
LOGGER.warn("Failed to collect pipe meta list from config node task agent", e);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,9 +25,11 @@
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStaticMeta;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStatus;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTemporaryMetaInCoordinator;
import org.apache.iotdb.confignode.consensus.response.pipe.task.PipeTableResp;
import org.apache.iotdb.confignode.manager.ConfigManager;
import org.apache.iotdb.confignode.manager.node.NodeManager;
import org.apache.iotdb.confignode.rpc.thrift.TShowPipeInfo;
import org.apache.iotdb.confignode.service.ConfigNode;
import org.apache.iotdb.rpc.TSStatusCode;

Expand Down Expand Up @@ -146,6 +148,25 @@ public void testFilter() {
Assert.assertEquals(3, allPipeTableResp.getAllPipeMeta().size());
}

@Test
public void testConvertToTShowPipeRespAggregatesRecentFailures() {
final PipeTableResp pipeTableResp = constructPipeTableResp();
final PipeTemporaryMetaInCoordinator temporaryMeta =
(PipeTemporaryMetaInCoordinator) pipeTableResp.getAllPipeMeta().get(0).getTemporaryMeta();
final Map<String, Long> firstNodeFailures = new HashMap<>();
firstNodeFailures.put("network_timeout", 10L);
firstNodeFailures.put("memory_timeout", 15L);
temporaryMeta.setRecentFailures(1, firstNodeFailures);
final Map<String, Long> secondNodeFailures = new HashMap<>();
secondNodeFailures.put("network_timeout", 2L);
temporaryMeta.setRecentFailures(2, secondNodeFailures);

final TShowPipeInfo showPipeInfo =
pipeTableResp.convertToTShowPipeResp().getPipeInfoList().get(0);
Assert.assertEquals(Long.valueOf(12), showPipeInfo.getRecentFailures().get("network_timeout"));
Assert.assertEquals(Long.valueOf(15), showPipeInfo.getRecentFailures().get("memory_timeout"));
}

@Test
public void testConvertToTShowPipeRespIncludesPreDeleteStatus() {
final PipeTableResp pipeTableResp = constructPipeTableResp();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStaticMeta;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStatus;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTemporaryMetaInCoordinator;
import org.apache.iotdb.confignode.consensus.request.write.pipe.task.CreatePipePlanV2;
import org.apache.iotdb.confignode.manager.ConfigManager;
import org.apache.iotdb.confignode.manager.ProcedureManager;
Expand Down Expand Up @@ -243,6 +244,50 @@ public void testParseHeartbeatDoesNotOverwritePreDeleteStatus() throws Exception
verify(context.procedureManager, never()).pipeHandleMetaChange(anyBoolean(), anyBoolean());
}

@Test
public void testParseHeartbeatAggregatesRecentFailuresFromAllDataNodes() throws Exception {
CommonDescriptor.getInstance().getConfig().setSeperatedPipeHeartbeatEnabled(false);

final String pipeName = "resourceFailurePipe";
final PipeTaskInfo pipeTaskInfo = new PipeTaskInfo();
createPipe(pipeTaskInfo, pipeName, PipeStatus.RUNNING);
final PipeMeta pipeMeta = pipeTaskInfo.getPipeMetaByPipeName(pipeName);
final ParserTestContext context = createParserTestContext(2, pipeTaskInfo);

final Map<String, Long> firstNodeFailures = new HashMap<>();
firstNodeFailures.put("network_timeout", 10L);
final Map<String, Long> secondNodeFailures = new HashMap<>();
secondNodeFailures.put("network_timeout", 2L);
secondNodeFailures.put("memory_timeout", 15L);

context.parser.parseHeartbeat(1, createPipeHeartbeat(pipeMeta, firstNodeFailures));
context.parser.parseHeartbeat(2, createPipeHeartbeat(pipeMeta, secondNodeFailures));

final PipeTemporaryMetaInCoordinator temporaryMeta =
(PipeTemporaryMetaInCoordinator) pipeMeta.getTemporaryMeta();
Assert.assertEquals(
Long.valueOf(12), temporaryMeta.getGlobalRecentFailures().get("network_timeout"));
Assert.assertEquals(
Long.valueOf(15), temporaryMeta.getGlobalRecentFailures().get("memory_timeout"));
verify(context.procedureManager, never()).pipeHandleMetaChange(anyBoolean(), anyBoolean());
}

@Test
public void testPipeHeartbeatTreatsNullRecentFailureMapAsEmpty() throws Exception {
final PipeTaskInfo pipeTaskInfo = new PipeTaskInfo();
createPipe(pipeTaskInfo, "nullFailureMapPipe", PipeStatus.RUNNING);
final PipeMeta pipeMeta = pipeTaskInfo.getPipeMetaByPipeName("nullFailureMapPipe");
final PipeHeartbeat heartbeat =
new PipeHeartbeat(
Collections.singletonList(pipeMeta.serialize()),
Collections.singletonList(false),
Collections.singletonList(0L),
Collections.singletonList(0d),
Collections.<Map<String, Long>>singletonList(null));

Assert.assertTrue(heartbeat.getRecentFailures(pipeMeta.getStaticMeta()).isEmpty());
}

private ParserTestContext createParserTestContext(final int registeredDataNodeCount) {
return createParserTestContext(registeredDataNodeCount, new PipeTaskInfo());
}
Expand Down Expand Up @@ -331,6 +376,16 @@ private PipeHeartbeat emptyHeartbeat() {
return new PipeHeartbeat(Collections.emptyList(), null, null, null);
}

private PipeHeartbeat createPipeHeartbeat(
final PipeMeta pipeMeta, final Map<String, Long> recentFailures) throws Exception {
return new PipeHeartbeat(
Collections.singletonList(pipeMeta.serialize()),
Collections.singletonList(false),
Collections.singletonList(0L),
Collections.singletonList(0d),
Collections.singletonList(recentFailures));
}

private static class ParserTestContext {
private final PipeHeartbeatParser parser;
private final ProcedureManager procedureManager;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStaticMeta;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStatus;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTemporaryMetaInAgent;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeType;
import org.apache.iotdb.commons.pipe.config.PipeConfig;
import org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant;
Expand Down Expand Up @@ -405,6 +406,7 @@ private void collectPipeMetaListInternal(final TDataNodeHeartbeatResp resp) thro
final List<Boolean> pipeCompletedList = new ArrayList<>();
final List<Long> pipeRemainingEventCountList = new ArrayList<>();
final List<Double> pipeRemainingTimeList = new ArrayList<>();
final List<Map<String, Long>> pipeRecentFailureList = new ArrayList<>();
try {
for (final PipeMeta pipeMeta : pipeMetaKeeper.getPipeMetaList()) {
pipeMetaBinaryList.add(pipeMeta.serialize());
Expand Down Expand Up @@ -438,6 +440,8 @@ private void collectPipeMetaListInternal(final TDataNodeHeartbeatResp resp) thro
pipeCompletedList.add(isCompleted);
pipeRemainingEventCountList.add(remainingEventAndTime.getLeft());
pipeRemainingTimeList.add(remainingEventAndTime.getRight());
pipeRecentFailureList.add(
((PipeTemporaryMetaInAgent) pipeMeta.getTemporaryMeta()).getRecentFailures());

logger.ifPresent(
l ->
Expand All @@ -457,6 +461,7 @@ private void collectPipeMetaListInternal(final TDataNodeHeartbeatResp resp) thro
resp.setPipeCompletedList(pipeCompletedList);
resp.setPipeRemainingEventCountList(pipeRemainingEventCountList);
resp.setPipeRemainingTimeList(pipeRemainingTimeList);
resp.setPipeRecentFailureList(pipeRecentFailureList);
PipeInsertionDataNodeListener.getInstance().listenToHeartbeat(true);
}

Expand Down Expand Up @@ -486,6 +491,7 @@ protected void collectPipeMetaListInternal(
final List<Boolean> pipeCompletedList = new ArrayList<>();
final List<Long> pipeRemainingEventCountList = new ArrayList<>();
final List<Double> pipeRemainingTimeList = new ArrayList<>();
final List<Map<String, Long>> pipeRecentFailureList = new ArrayList<>();
try {
for (final PipeMeta pipeMeta : pipeMetaKeeper.getPipeMetaList()) {
pipeMetaBinaryList.add(pipeMeta.serialize());
Expand Down Expand Up @@ -519,6 +525,8 @@ protected void collectPipeMetaListInternal(
pipeCompletedList.add(isCompleted);
pipeRemainingEventCountList.add(remainingEventAndTime.getLeft());
pipeRemainingTimeList.add(remainingEventAndTime.getRight());
pipeRecentFailureList.add(
((PipeTemporaryMetaInAgent) pipeMeta.getTemporaryMeta()).getRecentFailures());

logger.ifPresent(
l ->
Expand All @@ -538,6 +546,7 @@ protected void collectPipeMetaListInternal(
resp.setPipeCompletedList(pipeCompletedList);
resp.setPipeRemainingEventCountList(pipeRemainingEventCountList);
resp.setPipeRemainingTimeList(pipeRemainingTimeList);
resp.setPipeRecentFailureList(pipeRecentFailureList);
PipeInsertionDataNodeListener.getInstance().listenToHeartbeat(true);
}

Expand Down
Loading
Loading