Skip to content
Open
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 @@ -498,12 +498,50 @@ public TSchemaFetchResponse fetchSchema(TSchemaFetchRequest req) {
}

@Override
public TLoadResp sendTsFilePieceNode(TTsFilePieceReq req) {
LOGGER.info("Receive load node from uuid {}.", req.uuid);
public int getThriftMaxFrameSize() {
return IoTDBDescriptor.getInstance().getConfig().getThriftMaxFrameSize();
}

@Override
public TLoadResp sendTsFilePieceNode(final TTsFilePieceReq req) {
if (!req.isSetSliceIndex() || req.sliceIndex == 0) {
LOGGER.info("Receive load node from uuid {}.", req.uuid);
}

ConsensusGroupId groupId =
ConsensusGroupId.Factory.createFromTConsensusGroupId(req.consensusGroupId);
LoadTsFilePieceNode pieceNode = (LoadTsFilePieceNode) PlanNodeType.deserialize(req.body);
final boolean isSliced =
req.isSetSliceIndex() || req.isSetSliceCount() || req.isSetOriginBodySize();
if (isSliced) {
if (!req.isSetSliceIndex() || !req.isSetSliceCount() || !req.isSetOriginBodySize()) {
final List<String> missingFields = new ArrayList<>(3);
if (!req.isSetSliceIndex()) {
missingFields.add("sliceIndex");
}
if (!req.isSetSliceCount()) {
missingFields.add("sliceCount");
}
if (!req.isSetOriginBodySize()) {
missingFields.add("originBodySize");
}
return createTLoadResp(
RpcUtils.getStatus(
TSStatusCode.DESERIALIZE_PIECE_OF_TSFILE_ERROR,
String.format(
"Missing Load TsFile slice metadata: %s", String.join(", ", missingFields))));
}
return createTLoadResp(
StorageEngine.getInstance()
.writeLoadTsFileNodeSlice(
(DataRegionId) groupId,
req.body,
req.uuid,
req.sliceIndex,
req.sliceCount,
req.originBodySize));
}

final LoadTsFilePieceNode pieceNode = (LoadTsFilePieceNode) PlanNodeType.deserialize(req.body);
if (pieceNode == null) {
return createTLoadResp(
new TSStatus(TSStatusCode.DESERIALIZE_PIECE_OF_TSFILE_ERROR.getStatusCode()));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@

package org.apache.iotdb.db.queryengine.plan.scheduler.load;

import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId;
import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
import org.apache.iotdb.common.rpc.thrift.TEndPoint;
import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
Expand Down Expand Up @@ -51,17 +52,20 @@
import org.apache.iotdb.rpc.TSStatusCode;

import io.airlift.concurrent.SetThreadName;
import org.apache.thrift.TApplicationException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.io.IOException;
import java.net.SocketTimeoutException;
import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Future;
import java.util.concurrent.TimeoutException;
Expand All @@ -75,6 +79,7 @@ public class LoadTsFileDispatcherImpl implements IFragInstanceDispatcher, AutoCl

private static final int MAX_CONNECTION_TIMEOUT_MS = 24 * 60 * 60 * 1000; // 1 day
private static final int FIRST_ADJUSTMENT_TIMEOUT_MS = 6 * 60 * 60 * 1000; // 6 hours
private static final int LOAD_TSFILE_PIECE_RPC_FRAME_RESERVED_BYTES = 1024;
private static final AtomicInteger CONNECTION_TIMEOUT_MS =
new AtomicInteger(IoTDBDescriptor.getInstance().getConfig().getConnectionTimeoutInMS());

Expand All @@ -85,6 +90,7 @@ public class LoadTsFileDispatcherImpl implements IFragInstanceDispatcher, AutoCl
internalServiceClientManager;
private ExecutorService executor;
private final boolean isGeneratedByPipe;
private final Map<TEndPoint, Integer> endPoint2ThriftMaxFrameSize = new ConcurrentHashMap<>();

public LoadTsFileDispatcherImpl(
IClientManager<TEndPoint, SyncDataNodeInternalServiceClient> internalServiceClientManager,
Expand Down Expand Up @@ -134,26 +140,112 @@ public Future<FragInstanceDispatchResult> dispatch(

private void dispatchOneInstance(FragmentInstance instance)
throws FragmentInstanceDispatchException {
TTsFilePieceReq loadTsFileReq = null;
ByteBuffer body = null;

for (TDataNodeLocation dataNodeLocation :
instance.getRegionReplicaSet().getDataNodeLocations()) {
TEndPoint endPoint = dataNodeLocation.getInternalEndPoint();
if (isDispatchedToLocal(endPoint)) {
dispatchLocally(instance);
} else {
if (loadTsFileReq == null) {
loadTsFileReq =
new TTsFilePieceReq(
instance.getFragment().getPlanNodeTree().serializeToByteBuffer(),
uuid,
instance.getRegionReplicaSet().getRegionId());
if (body == null) {
body = instance.getFragment().getPlanNodeTree().serializeToByteBuffer();
}
dispatchRemote(loadTsFileReq, endPoint);
dispatchRemote(body, instance.getRegionReplicaSet().getRegionId(), endPoint);
}
}
}

private int getLoadTsFilePieceBodySizeLimit(final TEndPoint endPoint) throws Exception {
final int localMaxFrameSize = IoTDBDescriptor.getInstance().getConfig().getThriftMaxFrameSize();
Integer remoteMaxFrameSize = endPoint2ThriftMaxFrameSize.get(endPoint);
if (remoteMaxFrameSize == null) {
// An oversized frame is rejected before the RPC handler runs and closes the connection, so
// the receiver's frame limit cannot be recovered from the sender's transport exception.
try (SyncDataNodeInternalServiceClient client =
internalServiceClientManager.borrowClient(endPoint)) {
remoteMaxFrameSize = client.getThriftMaxFrameSize();
} catch (Exception e) {
if (!isUnknownMethod(e)) {
throw e;
}
// Older receivers still accept unsliced pieces. The failed RPC invalidates its client,
// so dispatchRemote borrows another client before sending the piece.
remoteMaxFrameSize = localMaxFrameSize;
}
if (remoteMaxFrameSize <= 0) {
throw new IllegalArgumentException(
String.format(
"Invalid Thrift maximum frame size %d from %s", remoteMaxFrameSize, endPoint));
}
endPoint2ThriftMaxFrameSize.put(endPoint, remoteMaxFrameSize);
}
return Math.max(
1,
Math.min(localMaxFrameSize, remoteMaxFrameSize)
- LOAD_TSFILE_PIECE_RPC_FRAME_RESERVED_BYTES);
}

private static boolean isUnknownMethod(Throwable e) {
do {
if (e instanceof TApplicationException
&& ((TApplicationException) e).getType() == TApplicationException.UNKNOWN_METHOD) {
return true;
}
} while ((e = e.getCause()) != null);
return false;
}

static List<TTsFilePieceReq> splitTsFilePieceReq(
final ByteBuffer body,
final String uuid,
final TConsensusGroupId consensusGroupId,
final int bodySizeLimit) {
if (bodySizeLimit <= 0) {
throw new IllegalArgumentException();
}

final int originBodySize = body.remaining();
final int sliceCount = getSliceCount(originBodySize, bodySizeLimit);
final List<TTsFilePieceReq> requests = new ArrayList<>(sliceCount);
if (sliceCount == 1) {
requests.add(createTsFilePieceReq(body.duplicate(), uuid, consensusGroupId));
return requests;
}

final int originPosition = body.position();
for (int sliceIndex = 0; sliceIndex < sliceCount; sliceIndex++) {
final int startOffset = sliceIndex * bodySizeLimit;
final int endOffset = startOffset + Math.min(bodySizeLimit, originBodySize - startOffset);
final ByteBuffer slicedBody = body.duplicate();
slicedBody.position(originPosition + startOffset);
slicedBody.limit(originPosition + endOffset);
requests.add(
createTsFilePieceReq(slicedBody.slice(), uuid, consensusGroupId)
.setSliceIndex(sliceIndex)
.setSliceCount(sliceCount)
.setOriginBodySize(originBodySize));
}
return requests;
}

static int getSliceCount(final int bodySize, final int bodySizeLimit) {
if (bodySize < 0 || bodySizeLimit <= 0) {
throw new IllegalArgumentException();
}
return bodySize == 0 ? 1 : (bodySize - 1) / bodySizeLimit + 1;
}

private static TTsFilePieceReq createTsFilePieceReq(
final ByteBuffer body, final String uuid, final TConsensusGroupId consensusGroupId) {
final TTsFilePieceReq request =
new TTsFilePieceReq().setUuid(uuid).setConsensusGroupId(consensusGroupId);
// The generated setter copies the whole buffer, while these immutable slices remain valid until
// all replicas have been dispatched.
request.body = body;
return request;
}

public void dispatchLocally(FragmentInstance instance) throws FragmentInstanceDispatchException {
if (isGeneratedByPipe) {
LOGGER.debug("Receive load node from uuid {}.", uuid);
Expand Down Expand Up @@ -211,24 +303,32 @@ public void dispatchLocally(FragmentInstance instance) throws FragmentInstanceDi
}
}

private void dispatchRemote(TTsFilePieceReq loadTsFileReq, TEndPoint endPoint)
private void dispatchRemote(
ByteBuffer body, TConsensusGroupId consensusGroupId, TEndPoint endPoint)
throws FragmentInstanceDispatchException {
try (SyncDataNodeInternalServiceClient client =
internalServiceClientManager.borrowClient(endPoint)) {
client.setTimeout(CONNECTION_TIMEOUT_MS.get());

final TLoadResp loadResp = client.sendTsFilePieceNode(loadTsFileReq);
if (!loadResp.isAccepted()) {
LOGGER.warn(loadResp.message);
throw new FragmentInstanceDispatchException(loadResp.status);
try {
final List<TTsFilePieceReq> loadTsFileReqs =
splitTsFilePieceReq(
body, uuid, consensusGroupId, getLoadTsFilePieceBodySizeLimit(endPoint));
try (SyncDataNodeInternalServiceClient client =
internalServiceClientManager.borrowClient(endPoint)) {
client.setTimeout(CONNECTION_TIMEOUT_MS.get());

for (final TTsFilePieceReq loadTsFileReq : loadTsFileReqs) {
final TLoadResp loadResp = client.sendTsFilePieceNode(loadTsFileReq);
if (!loadResp.isAccepted()) {
LOGGER.warn(loadResp.message);
throw new FragmentInstanceDispatchException(loadResp.status);
}
}
}
} catch (Exception e) {
adjustTimeoutIfNecessary(e);

final String exceptionMessage =
String.format(
"failed to dispatch load command %s to node %s because of exception: %s",
loadTsFileReq, endPoint, e);
uuid, endPoint, e);
LOGGER.warn(exceptionMessage, e);
throw new FragmentInstanceDispatchException(
new TSStatus()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@
import org.apache.iotdb.db.exception.load.LoadReadOnlyException;
import org.apache.iotdb.db.exception.runtime.StorageEngineFailureException;
import org.apache.iotdb.db.queryengine.plan.analyze.cache.schema.DataNodeTTLCache;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.PlanNodeType;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.load.LoadTsFilePieceNode;
import org.apache.iotdb.db.queryengine.plan.scheduler.load.LoadTsFileScheduler;
import org.apache.iotdb.db.service.metrics.FileMetrics;
Expand All @@ -73,6 +74,7 @@
import org.apache.iotdb.db.storageengine.dataregion.wal.exception.WALException;
import org.apache.iotdb.db.storageengine.dataregion.wal.recover.WALRecoverManager;
import org.apache.iotdb.db.storageengine.load.LoadTsFileManager;
import org.apache.iotdb.db.storageengine.load.LoadTsFilePieceNodeAssembler;
import org.apache.iotdb.db.storageengine.load.limiter.LoadTsFileRateLimiter;
import org.apache.iotdb.db.storageengine.rescon.memory.SystemInfo;
import org.apache.iotdb.db.utils.ThreadUtils;
Expand All @@ -88,6 +90,7 @@
import java.io.File;
import java.io.IOException;
import java.net.URL;
import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
Expand Down Expand Up @@ -976,6 +979,35 @@ public TSStatus writeLoadTsFileNode(
return RpcUtils.SUCCESS_STATUS;
}

public TSStatus writeLoadTsFileNodeSlice(
final DataRegionId dataRegionId,
final ByteBuffer body,
final String uuid,
final int sliceIndex,
final int sliceCount,
final int originBodySize) {
final LoadTsFilePieceNodeAssembler.Result result =
loadTsFileManager.appendPieceNodeSlice(
dataRegionId, uuid, body, sliceIndex, sliceCount, originBodySize);
if (!result.isValid()) {
return RpcUtils.getStatus(
TSStatusCode.DESERIALIZE_PIECE_OF_TSFILE_ERROR, result.getErrorMessage());
}
if (!result.isComplete()) {
return RpcUtils.SUCCESS_STATUS;
}

try {
final Object planNode = PlanNodeType.deserialize(result.getBody());
if (!(planNode instanceof LoadTsFilePieceNode)) {
return new TSStatus(TSStatusCode.DESERIALIZE_PIECE_OF_TSFILE_ERROR.getStatusCode());
}
return writeLoadTsFileNode(dataRegionId, (LoadTsFilePieceNode) planNode, uuid);
} catch (final Exception e) {
return new TSStatus(TSStatusCode.DESERIALIZE_PIECE_OF_TSFILE_ERROR.getStatusCode());
}
}

public TSStatus executeLoadCommand(
LoadTsFileScheduler.LoadCommand loadCommand,
String uuid,
Expand Down
Loading
Loading