Skip to content
Draft
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 @@ -1160,6 +1160,9 @@ public final class DataNodeQueryMessages {
"Start load TsFile {} locally.";
public static final String LOAD_ALL_FAILED_TSFILES_ARE_CONVERTED_TO_TABLETS =
"Load: all failed TsFiles are converted to tablets and inserted.";
public static final String
LOG_LOAD_FAILED_TO_LOAD_SOME_TSFILES_BY_CONVERTING_THEM_INTO_TABLETS_FAILED_TSFILES_ARG_7D9DB9C3 =
"Load: failed to load some TsFiles by converting them into tablets. Failed TsFiles: %s";

// --- Plan / Statement ---

Expand Down Expand Up @@ -2692,6 +2695,8 @@ public final class DataNodeQueryMessages {
"Parse or send TsFile %s error.";
public static final String DISPATCH_ONE_PIECE_TO_REPLICASET_ARG_ERROR_RESULT_STATUS_CODE_ARG =
"Dispatch one piece to ReplicaSet {} error. Result status code {}. ";
public static final String LOG_LOAD_CONSENSUS_SUBMIT_TRANSIENT_FAILURE_RETRY_D7E1D9A6 =
"Transient failure while submitting LOAD consensus {} (load {}) to {}, will retry ({}/{}): {}";
public static final String RESULT_STATUS_MESSAGE_ARG_DISPATCH_PIECE_NODE_ERROR_PERCENT_NARG =
"Result status message {}. Dispatch piece node error:%n{}";
public static final String SUB_STATUS_CODE_ARG_SUB_STATUS_MESSAGE_ARG =
Expand Down Expand Up @@ -3793,6 +3798,7 @@ private DataNodeQueryMessages() {}
public static final String EXCEPTION_THE_SECOND_ARGUMENT_OF_PERCENTILE_FUNCTION_PERCENTAGE_MUST_BE_A_DOUBLE_LITERAL_D9464B46 = "The second argument of 'percentile' function percentage must be a double literal";
public static final String EXCEPTION_DATA_TYPE_MISMATCH_FOR_MEASUREMENT_ARGARGARG_TYPE_IN_TSFILE_ARG_TYPE_IN_IOTDB_ARG_C5BA7DBD = "Data type mismatch for measurement %s%s%s, type in TsFile: %s, type in IoTDB: %s";
public static final String MESSAGE_FAILED_TO_RELEASE_EXTERNAL_TSFILE_QUERY_RESOURCE_712EE978 = "Failed to release external TsFile query resource";
public static final String EXCEPTION_UNKNOWN_LOADTSFILECONSENSUSOP_ORDINAL_ARG_62848FC2 = "Unknown LoadTsFileConsensusOp ordinal: ";
public static final String EXCEPTION_OUTER_QUERY_TIMEOUT_EXCEEDED_BEFORE_IOTDBLOCAL_QUERY_STARTS_800BFA63 = "Outer query timeout exceeded before IoTDBLocal query starts";
public static final String MESSAGE_FAILED_TO_CLOSE_UDF_RESULT_SET_AT_INDEX_ARG_A293B7EC = "Failed to close UDF result set at index {}";
public static final String EXCEPTION_INTERNAL_QUERY_EXECUTION_NOT_FOUND_62642542 = "Internal query execution not found";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,62 @@ private StorageEngineMessages() {}
// ======================== StorageEngine ========================

public static final String FAIL_TO_RECOVER_WAL = "Fail to recover wal.";
public static final String LOG_LOAD_CONSENSUS_WRITE_TO_REGION_ARG_VIA_PROTOCOL_ARG_EBB55042 =
"Write LOAD consensus node to region {} via protocol {}";
public static final String LOG_LOAD_CONSENSUS_WRITE_TO_REGION_ARG_VIA_PEER_ARG_FAILED_TRYING_NEXT_REPLICA_ARG_39217580 =
"LOAD consensus write to region {} via peer {} failed, trying next replica: {}";
public static final String LOG_LOAD_CONSENSUS_REFRESH_REPLICA_SET_FAILED_7C244C63 =
"Failed to refresh LOAD consensus replica set for region {}, using cached set: {}";
public static final String MESSAGE_LOAD_CONSENSUS_PIECE_CHECKSUM_MISMATCH_CF261675 =
"LOAD consensus piece checksum mismatch, loadId: %s, pieceIndex: %d";
public static final String MESSAGE_LOAD_CONSENSUS_WAL_FLUSH_FAILED_8BE1375A =
"Failed to flush LOAD consensus WAL entry";
public static final String MESSAGE_LOAD_CONSENSUS_RATIS_NOT_SUPPORTED_D371E344 =
"LOAD consensus is not supported on Ratis";
public static final String EXCEPTION_LOAD_CONSENSUS_STAGED_FILE_EOF_8743387D =
"Unexpected end of file when reading staged piece %s at offset %s.";
public static final String MESSAGE_LOAD_CONSENSUS_PIECE_NOT_CONTINUOUS_AFTER_FAILOVER_D6FFAC6C =
"LOAD piece %d of load %s cannot be applied because the previous pieces have not all been "
+ "applied (the staged state may have been lost by a leader failover).";
public static final String MESSAGE_LOAD_CONSENSUS_PIECE_DATA_MISSING_AFTER_PULL_8269CB0B =
"LOAD piece %d data of load %s is still missing after pulling from the write node.";
public static final String LOG_LOAD_CONSENSUS_FORWARD_PIECE_FAILED_34F9EBE7 =
"Failed to forward LOAD piece {} of load {} to follower {}: {}";
public static final String LOG_LOAD_CONSENSUS_PULL_PIECE_FAILED_AFB003D5 =
"Failed to pull LOAD piece {} of load {} from write node {}: {}";
public static final String MESSAGE_LOAD_CONSENSUS_PULL_WITHOUT_SOURCE_ENDPOINT_3B20D9E9 =
"LOAD pull request has no source endpoint.";
public static final String MESSAGE_LOAD_CONSENSUS_PULL_WITHOUT_RETAINED_PIECE_AD3C9D4F =
"LOAD piece %d of load %s is not retained on the write node.";
public static final String MESSAGE_LOAD_CONSENSUS_PULL_PUSH_BACK_FAILED_1A90C2B9 =
"Failed to push LOAD piece %d of load %s back to %s: %s";
public static final String LOG_LOAD_CONSENSUS_ABORT_MARKER_FAILED_6A218023 =
"Failed to log the LOAD ABORT marker of load {}: {}";
public static final String EXCEPTION_LOAD_CONSENSUS_PIECE_DATA_MISSING_OR_CHECKSUM_MISMATCH_AFTER_PULL_35F4972E =
"LOAD task %s piece %d data is missing or its checksum mismatches after pull.";
public static final String LOG_LOAD_CONSENSUS_RETAINED_PIECE_READ_FAILED_0659D19B =
"Failed to read retained LOAD piece {} of load {} from {}: {}";
public static final String LOG_LOAD_CONSENSUS_RETAINED_PIECE_WRITE_FAILED_99697608 =
"Failed to write retained LOAD piece {} of load {} to {}: {}";
public static final String LOG_LOAD_CONSENSUS_APPLIED_PIECE_RESTORE_FAILED_5BC74BBA =
"Failed to restore the applied LOAD piece entry {} of load {}.";
public static final String EXCEPTION_LOAD_CONSENSUS_STAGED_FILE_NOT_CONTINUOUS_F9408C19 =
"Staged file %s of load %s is not continuous: expected offset %d but current file length is %d.";
public static final String MESSAGE_LOAD_CONSENSUS_PREPARE_WITHOUT_STAGED_DATA_FE8ADC37 =
"Cannot prepare load %s because no staged data exists on this node.";
public static final String MESSAGE_LOAD_CONSENSUS_PREPARE_VERIFICATION_FAILED_B3865A82 =
"LOAD prepare verification failed for load %s: expected %d pieces with checksum %d, found %d pieces with checksum %d";
public static final String EXCEPTION_LOAD_CONSENSUS_STAGED_FILE_INCOMPLETE_1CDE954B =
"Staged file %s of load %s is incomplete and cannot be committed.";
public static final String LOG_LOAD_CONSENSUS_SNAPSHOT_TAKEN_09A7DD4C =
"Snapshotted %d in-progress LOAD task(s) with %d staged file(s) for region %s into %s.";
public static final String LOG_LOAD_CONSENSUS_SNAPSHOT_RESTORED_90ABC1BF =
"Restored %d in-progress LOAD task(s) with %d staged file(s) from snapshot %s.";
public static final String EXCEPTION_LOAD_CONSENSUS_SNAPSHOT_RESTORE_FAILED_F8C29C64 =
"Failed to restore LOAD snapshot from %s: %s";
public static final String EXCEPTION_LOAD_TSFILE_ALIGNED_VALUE_CHUNK_TIME_CHUNK_EEB00760 =
"Cannot attach value chunk of measurement %s in file %s: expected exactly one buffered "
+ "aligned time chunk, found %d.";
public static final String STORAGE_ENGINE_FAILED_TO_SET_UP = "Storage engine failed to set up.";
public static final String SEQ_MEMTABLE_FLUSH_CHECK_THREAD_STARTED = "start sequence memtable timed flush check thread successfully.";
public static final String UNSEQ_MEMTABLE_FLUSH_CHECK_THREAD_STARTED = "start unsequence memtable timed flush check thread successfully.";
Expand Down Expand Up @@ -477,6 +533,8 @@ private StorageEngineMessages() {}
public static final String CANNOT_CREATE_TSFILE_FOR_WRITING = "Can not create TsFile {} for writing.";
public static final String CLOSE_TSFILE_IO_WRITER_ERROR = "Close TsFileIOWriter {} error.";
public static final String CLOSE_MODIFICATION_FILE_ERROR = "Close ModificationFile {} error.";
public static final String EXCEPTION_TABLE_ARG_ARG_DOES_NOT_EXIST_WHEN_APPLYING_LOAD_CHUNK_DATA_IT_MAY_HAVE_BEEN_DROPPED_AFTER_THE_LOAD_WAS_ANALYZED_DDB35F93 =
"Table '%s.%s' does not exist when applying LOAD chunk data. It may have been dropped after the LOAD was analyzed.";
public static final String TASK_DIR_NOT_EMPTY_SKIP_DELETE = "Task dir {} is not empty, skip deleting.";
public static final String LOAD_CLEANUP_TASK_CANCELED = "Load cleanup task {} is canceled.";
public static final String LOAD_CLEANUP_TASK_STARTS = "Load cleanup task {} starts.";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1141,6 +1141,9 @@ public final class DataNodeQueryMessages {
"开始本地加载 TsFile {}。";
public static final String LOAD_ALL_FAILED_TSFILES_ARE_CONVERTED_TO_TABLETS =
"加载:所有失败的 TsFile 已转换为 Tablet 并插入。";
public static final String
LOG_LOAD_FAILED_TO_LOAD_SOME_TSFILES_BY_CONVERTING_THEM_INTO_TABLETS_FAILED_TSFILES_ARG_7D9DB9C3 =
"加载:部分 TsFile 通过转换为 Tablet 仍加载失败。失败的 TsFile:%s";

// --- Plan / Statement ---

Expand Down Expand Up @@ -3174,6 +3177,8 @@ public final class DataNodeQueryMessages {
public static final String DISPATCH_ONE_PIECE_TO_REPLICASET_ARG_ERROR_RESULT_STATUS_CODE_ARG =

"分发 TsFile 片段到 ReplicaSet {} 出错。结果状态码 {}。 ";
public static final String LOG_LOAD_CONSENSUS_SUBMIT_TRANSIENT_FAILURE_RETRY_D7E1D9A6 =
"提交 LOAD 共识 {}(load {})到 {} 时遇到瞬时失败,将重试({}/{}):{}";
public static final String RESULT_STATUS_MESSAGE_ARG_DISPATCH_PIECE_NODE_ERROR_PERCENT_NARG =

"结果状态消息 {}。分发片段节点出错:%n{}";
Expand Down Expand Up @@ -4549,6 +4554,7 @@ private DataNodeQueryMessages() {}
public static final String EXCEPTION_THE_SECOND_ARGUMENT_OF_PERCENTILE_FUNCTION_PERCENTAGE_MUST_BE_A_DOUBLE_LITERAL_D9464B46 = "'percentile' 函数的第二个参数 percentage 必须是 double 字面量";
public static final String EXCEPTION_DATA_TYPE_MISMATCH_FOR_MEASUREMENT_ARGARGARG_TYPE_IN_TSFILE_ARG_TYPE_IN_IOTDB_ARG_C5BA7DBD = "测点 %s%s%s 的数据类型不匹配,TsFile 中类型:%s,IoTDB 中类型:%s";
public static final String MESSAGE_FAILED_TO_RELEASE_EXTERNAL_TSFILE_QUERY_RESOURCE_712EE978 = "释放外部 TsFile 查询资源失败";
public static final String EXCEPTION_UNKNOWN_LOADTSFILECONSENSUSOP_ORDINAL_ARG_62848FC2 = "未知的 LoadTsFileConsensusOp 序号:";
public static final String EXCEPTION_OUTER_QUERY_TIMEOUT_EXCEEDED_BEFORE_IOTDBLOCAL_QUERY_STARTS_800BFA63 = "在 IoTDBLocal 查询开始前,外层查询已超时";
public static final String MESSAGE_FAILED_TO_CLOSE_UDF_RESULT_SET_AT_INDEX_ARG_A293B7EC = "关闭索引 {} 处的 UDF 结果集失败";
public static final String EXCEPTION_INTERNAL_QUERY_EXECUTION_NOT_FOUND_62642542 = "未找到内部查询执行";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,60 @@ private StorageEngineMessages() {}
// ======================== StorageEngine ========================

public static final String FAIL_TO_RECOVER_WAL = "WAL 恢复失败。";
public static final String LOG_LOAD_CONSENSUS_WRITE_TO_REGION_ARG_VIA_PROTOCOL_ARG_EBB55042 =
"通过协议 {} 向 Region {} 写入 LOAD 共识节点";
public static final String LOG_LOAD_CONSENSUS_WRITE_TO_REGION_ARG_VIA_PEER_ARG_FAILED_TRYING_NEXT_REPLICA_ARG_39217580 =
"向 Region {} 经节点 {} 写入 LOAD 共识失败,尝试下一个副本:{}";
public static final String LOG_LOAD_CONSENSUS_REFRESH_REPLICA_SET_FAILED_7C244C63 =
"刷新 Region {} 的 LOAD 共识副本集失败,使用缓存的副本集:{}";
public static final String MESSAGE_LOAD_CONSENSUS_PIECE_CHECKSUM_MISMATCH_CF261675 =
"LOAD 共识分片校验和不一致,loadId: %s,pieceIndex: %d";
public static final String MESSAGE_LOAD_CONSENSUS_WAL_FLUSH_FAILED_8BE1375A =
"LOAD 共识 WAL 记录落盘失败";
public static final String MESSAGE_LOAD_CONSENSUS_RATIS_NOT_SUPPORTED_D371E344 =
"Ratis 暂不支持 LOAD 共识";
public static final String EXCEPTION_LOAD_CONSENSUS_STAGED_FILE_EOF_8743387D =
"读取已暂存的分片 %s 时意外到达文件末尾,offset: %s。";
public static final String MESSAGE_LOAD_CONSENSUS_PIECE_NOT_CONTINUOUS_AFTER_FAILOVER_D6FFAC6C =
"LOAD 分片 %d(load %s)无法应用:之前的分片尚未全部应用(暂存状态可能在主备切换时丢失)。";
public static final String MESSAGE_LOAD_CONSENSUS_PIECE_DATA_MISSING_AFTER_PULL_8269CB0B =
"LOAD 分片 %d(load %s)的数据在向写节点回补后仍然缺失。";
public static final String LOG_LOAD_CONSENSUS_FORWARD_PIECE_FAILED_34F9EBE7 =
"向副本 {} 转发 LOAD 分片 {}(load {})失败:{}";
public static final String LOG_LOAD_CONSENSUS_PULL_PIECE_FAILED_AFB003D5 =
"向写节点 {} 回补 LOAD 分片 {}(load {})失败:{}";
public static final String MESSAGE_LOAD_CONSENSUS_PULL_WITHOUT_SOURCE_ENDPOINT_3B20D9E9 =
"LOAD 回补请求缺少源节点地址。";
public static final String MESSAGE_LOAD_CONSENSUS_PULL_WITHOUT_RETAINED_PIECE_AD3C9D4F =
"LOAD 分片 %d(load %s)未在写节点保留。";
public static final String MESSAGE_LOAD_CONSENSUS_PULL_PUSH_BACK_FAILED_1A90C2B9 =
"将 LOAD 分片 %d(load %s)推回给 %s 失败:%s";
public static final String LOG_LOAD_CONSENSUS_ABORT_MARKER_FAILED_6A218023 =
"写入 LOAD 中止(ABORT)标记(load {})失败:{}";
public static final String EXCEPTION_LOAD_CONSENSUS_PIECE_DATA_MISSING_OR_CHECKSUM_MISMATCH_AFTER_PULL_35F4972E =
"回补后 LOAD 任务 %s 的分片 %d 数据仍缺失或校验和不一致。";
public static final String LOG_LOAD_CONSENSUS_RETAINED_PIECE_READ_FAILED_0659D19B =
"读取保留的 LOAD 分片 {}(load {})失败,文件:{},原因:{}";
public static final String LOG_LOAD_CONSENSUS_RETAINED_PIECE_WRITE_FAILED_99697608 =
"写入保留的 LOAD 分片 {}(load {})失败,文件:{},原因:{}";
public static final String LOG_LOAD_CONSENSUS_APPLIED_PIECE_RESTORE_FAILED_5BC74BBA =
"恢复 load {} 的已应用 LOAD 分片条目 {} 失败。";
public static final String EXCEPTION_LOAD_CONSENSUS_STAGED_FILE_NOT_CONTINUOUS_F9408C19 =
"load %s 的暂存文件 %s 不连续:期望偏移 %d,但当前文件长度为 %d。";
public static final String MESSAGE_LOAD_CONSENSUS_PREPARE_WITHOUT_STAGED_DATA_FE8ADC37 =
"无法准备(PREPARE)load %s,因为该节点上不存在暂存数据。";
public static final String MESSAGE_LOAD_CONSENSUS_PREPARE_VERIFICATION_FAILED_B3865A82 =
"LOAD PREPARE 校验失败,load %s:预期 %d 片、checksum %d,实际 %d 片、checksum %d";
public static final String EXCEPTION_LOAD_CONSENSUS_STAGED_FILE_INCOMPLETE_1CDE954B =
"load %s 的暂存文件 %s 不完整,无法提交(COMMIT)。";
public static final String LOG_LOAD_CONSENSUS_SNAPSHOT_TAKEN_09A7DD4C =
"已将 region %s 的 %d 个进行中的 LOAD 任务(共 %d 个暂存文件)纳入快照 %s。";
public static final String LOG_LOAD_CONSENSUS_SNAPSHOT_RESTORED_90ABC1BF =
"已从快照 %s 恢复 %d 个进行中的 LOAD 任务(共 %d 个暂存文件)。";
public static final String EXCEPTION_LOAD_CONSENSUS_SNAPSHOT_RESTORE_FAILED_F8C29C64 =
"从 %s 恢复 LOAD 快照失败:%s";
public static final String EXCEPTION_LOAD_TSFILE_ALIGNED_VALUE_CHUNK_TIME_CHUNK_EEB00760 =
"无法将测量 %s 的值 Chunk 挂载到文件 %s:预期恰好一个已缓冲的 Aligned 时间 Chunk,实际发现 %d 个。";
public static final String STORAGE_ENGINE_FAILED_TO_SET_UP = "存储引擎启动失败。";
public static final String SEQ_MEMTABLE_FLUSH_CHECK_THREAD_STARTED = "顺序 memtable 定时 flush 检查线程启动成功。";
public static final String UNSEQ_MEMTABLE_FLUSH_CHECK_THREAD_STARTED = "乱序 memtable 定时 flush 检查线程启动成功。";
Expand Down Expand Up @@ -477,6 +531,8 @@ private StorageEngineMessages() {}
public static final String CANNOT_CREATE_TSFILE_FOR_WRITING = "无法创建 TsFile {} 用于写入。";
public static final String CLOSE_TSFILE_IO_WRITER_ERROR = "关闭 TsFileIOWriter {} 出错。";
public static final String CLOSE_MODIFICATION_FILE_ERROR = "关闭修改文件 {} 出错。";
public static final String EXCEPTION_TABLE_ARG_ARG_DOES_NOT_EXIST_WHEN_APPLYING_LOAD_CHUNK_DATA_IT_MAY_HAVE_BEEN_DROPPED_AFTER_THE_LOAD_WAS_ANALYZED_DDB35F93 =
"应用 LOAD chunk 数据时表 '%s.%s' 不存在,可能在 LOAD 分析之后被删除了。";
public static final String TASK_DIR_NOT_EMPTY_SKIP_DELETE = "任务目录 {} 非空,跳过删除。";
public static final String LOAD_CLEANUP_TASK_CANCELED = "加载清理任务 {} 已取消。";
public static final String LOAD_CLEANUP_TASK_STARTS = "加载清理任务 {} 开始。";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4351,7 +4351,7 @@ public int getLoadTsFileSpiltPartitionMaxSize() {
}

public void setLoadTsFileSpiltPartitionMaxSize(int loadTsFileSpiltPartitionMaxSize) {
if (loadTsFileSpiltPartitionMaxSize <= 0) {
if (loadTsFileSpiltPartitionMaxSize < 0) {
throw new IllegalArgumentException(
DataNodeMiscMessages
.MISC_EXCEPTION_LOADTSFILESPILTPARTITIONMAXSIZE_SHOULD_BE_GREATER_THAN_OR_95B4DB23);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@
import org.apache.iotdb.db.i18n.DataNodeMiscMessages;
import org.apache.iotdb.db.i18n.DataNodePipeMessages;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.PlanVisitor;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.load.LoadTsFileConsensusNode;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.pipe.PipeEnrichedDeleteDataNode;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.pipe.PipeEnrichedInsertNode;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.DeleteDataNode;
Expand All @@ -47,6 +48,7 @@
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.RelationalInsertRowNode;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.RelationalInsertRowsNode;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.RelationalInsertTabletNode;
import org.apache.iotdb.db.storageengine.StorageEngine;
import org.apache.iotdb.db.storageengine.dataregion.DataRegion;
import org.apache.iotdb.rpc.RpcUtils;
import org.apache.iotdb.rpc.TSStatusCode;
Expand Down Expand Up @@ -316,4 +318,16 @@ public TSStatus visitPipeEnrichedDeleteDataNode(
public TSStatus visitWriteObjectFile(ObjectNode node, DataRegion dataRegion) {
throw new UnsupportedOperationException();
}

@Override
public TSStatus visitLoadTsFileConsensus(LoadTsFileConsensusNode node, DataRegion dataRegion) {
try {
return StorageEngine.getInstance()
.getLoadTsFileManager()
.applyConsensusRequest(dataRegion, node);
} catch (Exception e) {
LOGGER.error(DataNodeMiscMessages.ERROR_EXECUTING_PLAN_NODE, node, e);
return new TSStatus(TSStatusCode.LOAD_FILE_ERROR.getStatusCode()).setMessage(e.getMessage());
}
}
}
Loading
Loading