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 @@ -23,6 +23,9 @@
public final class ImportWALMessages {

public static final String MESSAGE_IMPORT_WAL_5E42804E = "import-wal";
public static final String
MESSAGE_UNSUPPORTED_WALS_1_WALS_WRITTEN_WITH_THE_ELASTICSTRATEGY_2_WALS_WRITTEN_WITH_THE_ROUNDROBINSTRATEGY_REPLAYING_THEM_MAY_FAIL_C542F812 =
"Unsupported WALs: (1) WALs written with the ElasticStrategy; (2) WALs written with the RoundRobinStrategy. Replaying them may fail.";
public static final String
MESSAGE_PATH_OF_A_WAL_FILE_OR_A_DIRECTORY_CONTAINING_WAL_FILES_473D0554 =
"Path of a WAL file or a directory containing WAL files.";
Expand Down Expand Up @@ -80,6 +83,9 @@ public final class ImportWALMessages {
"Unsupported on_success value: %s. Expected none or delete.";
public static final String EXCEPTION_TABLE_MODEL_WAL_ENTRIES_REQUIRE_DB_DATABASE_F7597726 =
"Table-model WAL entries require -db/--database.";
public static final String
EXCEPTION_A_WAL_SNAPSHOT_REQUIRES_A_DECLARED_TARGET_DATABASE_TO_DETERMINE_ITS_DATA_MODEL_SPECIFY_DB_DATABASE_382FC74C =
"A WAL snapshot requires a declared target database to determine its data model. Specify -db/--database.";
public static final String EXCEPTION_UNSUPPORTED_WAL_OPERATION_ARG_ABD227A0 =
"Unsupported WAL operation: %s";
public static final String
Expand Down Expand Up @@ -136,14 +142,16 @@ public final class ImportWALMessages {
public static final String MESSAGE_SKIPPED_ARG_CORRUPTED_WAL_FILES_SOURCE_FILES_RETAINED_A889CCE2 =
"Skipped %d corrupted WAL files; source files retained.";

public static final String MESSAGE_TARGET_DATABASE_FOR_TABLE_MODEL_WAL_ENTRIES_IF_OMITTED_INFER_FROM_THE_WAL_PARENT_DIRECTORY_AND_ASK_FOR_CONFIRMATION_4B1E409D =
"Target database for table-model WAL entries. If omitted, infer from the WAL parent directory and ask for confirmation.";
public static final String MESSAGE_TARGET_DATABASE_FOR_WAL_REPLAY_IF_OMITTED_INFER_FROM_THE_WAL_PARENT_DIRECTORY_IF_INFERENCE_FAILS_DB_DATABASE_IS_REQUIRED_6EEBC019 =
"Target database for WAL replay. If omitted, infer from the WAL parent directory; if inference fails, -db/--database is required.";
public static final String MESSAGE_INFERRED_TABLE_DATABASE_ARG_FROM_WAL_DIRECTORY_ARG_REPLAY_INTO_THIS_DATABASE_Y_YES_A_ACCEPT_ALL_INFERRED_DATABASES_N_QUIT_5B59D833 =
"Inferred table database %s from WAL directory %s. Replay into this database? [y] yes, [a] accept all inferred databases, [N] quit: ";
public static final String EXCEPTION_DATABASE_CONFIRMATION_REQUIRED_FOR_WAL_DIRECTORY_ARG_INFERRED_DATABASE_ARG_SPECIFY_DB_DATABASE_OR_SKIP_DB_CONFIRMATION_WHEN_INTERACTIVE_INPUT_IS_UNAVAILABLE_14DF6D36 =
"Database confirmation required for WAL directory %s (inferred database: %s). Specify -db/--database or --skip_db_confirmation when interactive input is unavailable.";
public static final String EXCEPTION_REPLAY_INTO_INFERRED_DATABASE_ARG_WAS_NOT_CONFIRMED_SPECIFY_DB_DATABASE_TO_SELECT_THE_TARGET_EXPLICITLY_86F81190 =
"Replay into inferred database %s was not confirmed. Specify -db/--database to select the target explicitly.";
public static final String EXCEPTION_CANNOT_DETERMINE_THE_TARGET_DATABASE_OF_WAL_DIRECTORIES_ARG_SPECIFY_DB_DATABASE_WHICH_APPLIES_TO_ALL_IMPORTED_DIRECTORIES_55B174E4 =
"Cannot determine the target database of WAL directories %s. Specify -db/--database, which applies to all imported directories.";

public static final String MESSAGE_ACCEPT_ALL_INFERRED_DATABASE_NAMES_WITHOUT_CONFIRMATION_DB_DATABASE_STILL_TAKES_PRECEDENCE_FA49A73C =
"Accept all inferred database names without confirmation; -db/--database still takes precedence.";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,9 @@
public final class ImportWALMessages {

public static final String MESSAGE_IMPORT_WAL_5E42804E = "import-wal";
public static final String
MESSAGE_UNSUPPORTED_WALS_1_WALS_WRITTEN_WITH_THE_ELASTICSTRATEGY_2_WALS_WRITTEN_WITH_THE_ROUNDROBINSTRATEGY_REPLAYING_THEM_MAY_FAIL_C542F812 =
"不支持导入以下两类 WAL:(1) 使用 ElasticStrategy 写入的 WAL;(2) 使用 RoundRobinStrategy 写入的 WAL。强行导入可能会出错。";
public static final String
MESSAGE_PATH_OF_A_WAL_FILE_OR_A_DIRECTORY_CONTAINING_WAL_FILES_473D0554 =
"WAL 文件或包含 WAL 文件的目录路径。";
Expand Down Expand Up @@ -79,6 +82,9 @@ public final class ImportWALMessages {
"不支持的 on_success 值:%s。应为 none 或 delete。";
public static final String EXCEPTION_TABLE_MODEL_WAL_ENTRIES_REQUIRE_DB_DATABASE_F7597726 =
"表模型 WAL 条目要求指定 -db/--database。";
public static final String
EXCEPTION_A_WAL_SNAPSHOT_REQUIRES_A_DECLARED_TARGET_DATABASE_TO_DETERMINE_ITS_DATA_MODEL_SPECIFY_DB_DATABASE_382FC74C =
"WAL 快照需要先声明目标数据库才能确定其数据模型。请指定 -db/--database。";
public static final String EXCEPTION_UNSUPPORTED_WAL_OPERATION_ARG_ABD227A0 =
"不支持的 WAL 操作:%s";
public static final String
Expand Down Expand Up @@ -135,14 +141,16 @@ public final class ImportWALMessages {
public static final String MESSAGE_SKIPPED_ARG_CORRUPTED_WAL_FILES_SOURCE_FILES_RETAINED_A889CCE2 =
"已跳过 %d 个损坏的 WAL 文件;源文件已保留。";

public static final String MESSAGE_TARGET_DATABASE_FOR_TABLE_MODEL_WAL_ENTRIES_IF_OMITTED_INFER_FROM_THE_WAL_PARENT_DIRECTORY_AND_ASK_FOR_CONFIRMATION_4B1E409D =
"表模型 WAL 条目的目标数据库。省略时尝试从 WAL 文件的父目录名推断,并请求确认。";
public static final String MESSAGE_TARGET_DATABASE_FOR_WAL_REPLAY_IF_OMITTED_INFER_FROM_THE_WAL_PARENT_DIRECTORY_IF_INFERENCE_FAILS_DB_DATABASE_IS_REQUIRED_6EEBC019 =
"WAL 重放的目标数据库。省略时尝试从 WAL 文件的父目录名推断;无法推断时必须指定 -db/--database。";
public static final String MESSAGE_INFERRED_TABLE_DATABASE_ARG_FROM_WAL_DIRECTORY_ARG_REPLAY_INTO_THIS_DATABASE_Y_YES_A_ACCEPT_ALL_INFERRED_DATABASES_N_QUIT_5B59D833 =
"推断目标表模型数据库为 %s,来源 WAL 目录为 %s。是否向此数据库重放?[y] 同意,[a] 全部同意推断的数据库,[N] 退出:";
public static final String EXCEPTION_DATABASE_CONFIRMATION_REQUIRED_FOR_WAL_DIRECTORY_ARG_INFERRED_DATABASE_ARG_SPECIFY_DB_DATABASE_OR_SKIP_DB_CONFIRMATION_WHEN_INTERACTIVE_INPUT_IS_UNAVAILABLE_14DF6D36 =
"需要确认 WAL 目录 %s 的目标数据库(推断结果:%s)。无交互终端时请指定 -db/--database 或 --skip_db_confirmation。";
public static final String EXCEPTION_REPLAY_INTO_INFERRED_DATABASE_ARG_WAS_NOT_CONFIRMED_SPECIFY_DB_DATABASE_TO_SELECT_THE_TARGET_EXPLICITLY_86F81190 =
"未确认向推断出的数据库 %s 重放。请使用 -db/--database 显式选择目标。";
public static final String EXCEPTION_CANNOT_DETERMINE_THE_TARGET_DATABASE_OF_WAL_DIRECTORIES_ARG_SPECIFY_DB_DATABASE_WHICH_APPLIES_TO_ALL_IMPORTED_DIRECTORIES_55B174E4 =
"无法推断 WAL 目录 %s 的目标数据库,请使用 -db/--database 显式指定(-db 会作用于本次导入的所有目录)。";

public static final String MESSAGE_ACCEPT_ALL_INFERRED_DATABASE_NAMES_WITHOUT_CONFIRMATION_DB_DATABASE_STILL_TAKES_PRECEDENCE_FA49A73C =
"自动接受所有推断出的数据库名,不再询问确认;-db/--database 仍优先。";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -200,9 +200,10 @@ public PipeRawTabletInsertionEvent(
}

@TestOnly
public PipeRawTabletInsertionEvent(final Tablet tablet, final boolean isAligned) {
public PipeRawTabletInsertionEvent(
final boolean isTableModel, final Tablet tablet, final boolean isAligned) {
this(
null,
isTableModel,
null,
null,
null,
Expand All @@ -225,9 +226,12 @@ public PipeRawTabletInsertionEvent(final Tablet tablet, final boolean isAligned)

@TestOnly
public PipeRawTabletInsertionEvent(
final Tablet tablet, final boolean isAligned, final TreePattern treePattern) {
final boolean isTableModel,
final Tablet tablet,
final boolean isAligned,
final TreePattern treePattern) {
this(
null,
isTableModel,
null,
null,
null,
Expand All @@ -250,10 +254,27 @@ public PipeRawTabletInsertionEvent(

@TestOnly
public PipeRawTabletInsertionEvent(
final Tablet tablet, final long startTime, final long endTime) {
final boolean isTableModel, final Tablet tablet, final long startTime, final long endTime) {
this(
null, null, null, null, tablet, false, null, false, null, 0, null, null, null, null, null,
null, true, startTime, endTime);
isTableModel,
null,
null,
null,
tablet,
false,
null,
false,
null,
0,
null,
null,
null,
null,
null,
null,
true,
startTime,
endTime);
}

@Override
Expand Down Expand Up @@ -479,8 +500,10 @@ public Tablet convertToTablet() {

private TabletInsertionEventParser initEventParser() {
if (eventParser == null) {
// The data model of the payload is the one declared by the event, it must never be derived
// from the device name of the tablet.
eventParser =
tablet.getDeviceId().startsWith("root.")
!isTableModelEvent()
? new TabletInsertionEventTreePatternParser(
pipeTaskMeta,
this,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -320,7 +320,8 @@ public TsFileInsertionEventQueryParser(
measurement,
meta.getStatistics().getStartTime(),
meta.getStatistics().getEndTime(),
currentModifications);
currentModifications,
false);
} catch (IOException e) {
LOGGER.warn(
DataNodePipeMessages.FAILED_TO_READ_METADATA_FOR_DEVICEID_MEASUREMENT,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -115,7 +115,7 @@ public class TsFileInsertionEventQueryParserTabletIterator implements Iterator<T

this.measurementModsList =
ModsOperationUtil.initializeMeasurementMods(
deviceId, this.measurements, currentModifications);
deviceId, this.measurements, currentModifications, false);
}

private QueryDataSet buildQueryDataSet() throws IOException {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -691,7 +691,10 @@ private void moveToNextChunkReader()
chunkHeader.getCompressionType()));
modsInfos.addAll(
ModsOperationUtil.initializeMeasurementMods(
currentDevice, Collections.singletonList(measurementID), currentModifications));
currentDevice,
Collections.singletonList(measurementID),
currentModifications,
false));
return;
}
case MetaMarker.VALUE_CHUNK_HEADER:
Expand Down Expand Up @@ -830,7 +833,8 @@ private boolean filterChunk(
chunkHeader.getMeasurementID(),
statistics.getStartTime(),
statistics.getEndTime(),
currentModifications)) {
currentModifications,
false)) {
tsFileSequenceReader.position(nextMarkerOffset);
return true;
}
Expand Down Expand Up @@ -999,7 +1003,7 @@ private void cacheAlignedValueChunk(final CachedAlignedValueChunk valueChunk) th
chunkHeader.getCompressionType()));
pendingAlignedChunkGroup.modsInfos.addAll(
ModsOperationUtil.initializeMeasurementMods(
currentDevice, Collections.singletonList(measurementID), currentModifications));
currentDevice, Collections.singletonList(measurementID), currentModifications, false));
}

private PendingAlignedChunkGroup getOrCreatePendingAlignedChunkGroup(final int timeChunkIndex) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -466,7 +466,7 @@ private void initChunkReader(final AbstractAlignedChunkMetadata alignedChunkMeta

this.chunkReader = new TableChunkReader(timeChunk, valueChunkList, null);
this.modsInfoList =
ModsOperationUtil.initializeMeasurementMods(deviceID, measurementList, modifications);
ModsOperationUtil.initializeMeasurementMods(deviceID, measurementList, modifications, true);
}

private boolean areAllFieldsDeletedByMods(
Expand All @@ -481,7 +481,8 @@ private boolean areAllFieldsDeletedByMods(
internMeasurementName(schema),
alignedChunkMetadata.getStartTime(),
alignedChunkMetadata.getEndTime(),
modifications)) {
modifications,
true)) {
return false;
}
}
Expand All @@ -492,7 +493,7 @@ private boolean isFieldDeletedByMods(
final String measurementID, final long startTime, final long endTime) {
return !modifications.isEmpty()
&& ModsOperationUtil.isAllDeletedByMods(
deviceID, measurementID, startTime, endTime, modifications);
deviceID, measurementID, startTime, endTime, modifications, true);
}

private String internMeasurementName(final IMeasurementSchema schema) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -80,25 +80,28 @@ private ModsOperationUtil() {
* @param startTime start time
* @param endTime end time
* @param modifications modification records
* @param isTableModel whether the device belongs to table model
* @return true if data is completely deleted, false otherwise
*/
public static boolean isAllDeletedByMods(
IDeviceID deviceID,
String measurementID,
long startTime,
long endTime,
PatternTreeMap<ModEntry, PatternTreeMapFactory.ModsSerializer> modifications) {
PatternTreeMap<ModEntry, PatternTreeMapFactory.ModsSerializer> modifications,
boolean isTableModel) {
if (modifications == null) {
return false;
}

final List<ModEntry> mods = getOverlappedMods(deviceID, measurementID, modifications);
final List<ModEntry> mods =
getOverlappedMods(deviceID, measurementID, modifications, isTableModel);
if (mods == null || mods.isEmpty()) {
return false;
}

// Different logic for tree model and table model
if (deviceID.isTableModel()) {
if (isTableModel) {
// For table model: check if any modification affects the device and covers the time range
return mods.stream()
.anyMatch(
Expand All @@ -119,17 +122,20 @@ public static boolean isAllDeletedByMods(
* @param deviceID device ID
* @param measurements measurement list
* @param modifications modification records
* @param isTableModel whether the device belongs to table model
* @return mapping from measurement ID to mods list and index
*/
public static List<ModsInfo> initializeMeasurementMods(
IDeviceID deviceID,
List<String> measurements,
PatternTreeMap<ModEntry, PatternTreeMapFactory.ModsSerializer> modifications) {
PatternTreeMap<ModEntry, PatternTreeMapFactory.ModsSerializer> modifications,
boolean isTableModel) {

List<ModsInfo> modsInfos = new ArrayList<>(measurements.size());

for (final String measurement : measurements) {
final List<ModEntry> mods = getOverlappedMods(deviceID, measurement, modifications);
final List<ModEntry> mods =
getOverlappedMods(deviceID, measurement, modifications, isTableModel);
if (mods == null || mods.isEmpty()) {
// No mods, use empty list and index 0
modsInfos.add(new ModsInfo(Collections.emptyList(), 0));
Expand All @@ -139,7 +145,7 @@ public static List<ModsInfo> initializeMeasurementMods(
// Sort by time range for efficient lookup
// Different filtering logic for tree model and table model
final List<ModEntry> filteredMods;
if (deviceID.isTableModel()) {
if (isTableModel) {
// For table model: filter modifications that affect the device
filteredMods =
mods.stream()
Expand All @@ -161,9 +167,11 @@ public static List<ModsInfo> initializeMeasurementMods(
private static List<ModEntry> getOverlappedMods(
final IDeviceID deviceID,
final String measurement,
final PatternTreeMap<ModEntry, PatternTreeMapFactory.ModsSerializer> modifications) {
final PatternTreeMap<ModEntry, PatternTreeMapFactory.ModsSerializer> modifications,
final boolean isTableModel) {
try {
return modifications.getOverlapped(CompactionPathUtils.getPath(deviceID, measurement));
return modifications.getOverlapped(
CompactionPathUtils.getPath(deviceID, measurement, isTableModel));
} catch (final IllegalPathException e) {
throw new PipeException(e.getMessage(), e);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -329,7 +329,8 @@ private void writeTabletsIntoOneFile(
}
}

final MemTableFlushTask memTableFlushTask = new MemTableFlushTask(memTable, writer, null, null);
final MemTableFlushTask memTableFlushTask =
new MemTableFlushTask(memTable, writer, null, null, true);
memTableFlushTask.syncFlushMemTable();

writer.endFile();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -257,7 +257,8 @@ private void writeTabletsIntoOneFile(
}
}

final MemTableFlushTask memTableFlushTask = new MemTableFlushTask(memTable, writer, null, null);
final MemTableFlushTask memTableFlushTask =
new MemTableFlushTask(memTable, writer, null, null, false);
memTableFlushTask.syncFlushMemTable();

writer.endFile();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -71,7 +71,6 @@
import org.apache.iotdb.pipe.api.exception.PipeParameterNotValidException;

import org.apache.tsfile.file.metadata.IDeviceID;
import org.apache.tsfile.file.metadata.PlainDeviceID;
import org.apache.tsfile.utils.Pair;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
Expand Down Expand Up @@ -124,17 +123,13 @@
import static org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_START_TIME_KEY;
import static org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.SOURCE_TSFILE_PARSER_KEY;
import static org.apache.iotdb.commons.pipe.source.IoTDBSource.getSkipIfNoPrivileges;
import static org.apache.tsfile.common.constant.TsFileConstant.PATH_ROOT;
import static org.apache.tsfile.common.constant.TsFileConstant.PATH_SEPARATOR;

public class PipeHistoricalDataRegionTsFileAndDeletionSource
implements PipeHistoricalDataRegionSource {

private static final Logger LOGGER =
LoggerFactory.getLogger(PipeHistoricalDataRegionTsFileAndDeletionSource.class);

private static final String TREE_MODEL_EVENT_TABLE_NAME_PREFIX = PATH_ROOT + PATH_SEPARATOR;

private String pipeName;
private long creationTime;
private String pipeNameWithCreationTime;
Expand Down Expand Up @@ -1026,7 +1021,7 @@ private boolean mayTsFileResourceOverlappedWithPattern(final TsFileResource reso
.anyMatch(
deviceID -> {
if (!isModelDetected) {
detectModel(resource, deviceID);
detectModel(resource);
isModelDetected = true;
}

Expand All @@ -1039,13 +1034,10 @@ private boolean mayTsFileResourceOverlappedWithPattern(final TsFileResource reso
});
}

private void detectModel(final TsFileResource resource, final IDeviceID deviceID) {
this.isTableModel =
!(deviceID instanceof PlainDeviceID
|| deviceID.getTableName().startsWith(TREE_MODEL_EVENT_TABLE_NAME_PREFIX)
|| deviceID.getTableName().equals(PATH_ROOT));

private void detectModel(final TsFileResource resource) {
// One source serves one data region, so the data model is decided by its database.
final String databaseName = resource.getDatabaseName();
this.isTableModel = PathUtils.isTableModelDatabase(databaseName);
isDbNameCoveredByPattern =
isTableModel
? tablePattern.isTableModelDataAllowedToBeCaptured()
Expand Down
Loading
Loading