From 9519e39510217a2288984d4afec05e9543765459 Mon Sep 17 00:00:00 2001 From: wangxiang Date: Fri, 11 Sep 2026 15:59:44 +0800 Subject: [PATCH] Improve comments in tsfile. --- .../java/org/apache/tsfile/enums/ColumnCategory.java | 4 ++++ .../org/apache/tsfile/encoding/encoder/SDTEncoder.java | 5 +++++ .../java/org/apache/tsfile/file/header/PageHeader.java | 9 ++++++++- .../apache/tsfile/file/metadata/MetadataIndexNode.java | 8 ++++++++ .../tsfile/file/metadata/TimeseriesMetadata.java | 4 ++++ .../apache/tsfile/file/metadata/TsFileMetadata.java | 7 +++++++ .../org/apache/tsfile/read/TsFileSequenceReader.java | 7 +++++++ .../main/java/org/apache/tsfile/read/common/Chunk.java | 10 ++++++++++ .../org/apache/tsfile/read/reader/IChunkReader.java | 4 ++++ .../org/apache/tsfile/read/reader/IPageReader.java | 4 ++++ .../apache/tsfile/read/reader/chunk/ChunkReader.java | 9 +++++++++ .../org/apache/tsfile/read/reader/page/PageReader.java | 4 ++++ .../apache/tsfile/read/v4/DeviceTableModelReader.java | 5 +++++ .../java/org/apache/tsfile/write/TsFileWriter.java | 6 +++++- .../write/chunk/AlignedChunkGroupWriterImpl.java | 9 +++++++++ .../org/apache/tsfile/write/chunk/ChunkWriterImpl.java | 5 +++++ .../apache/tsfile/write/chunk/IChunkGroupWriter.java | 6 ++++++ .../apache/tsfile/write/schema/IMeasurementSchema.java | 4 ++++ .../apache/tsfile/write/schema/MeasurementSchema.java | 5 ++++- .../java/org/apache/tsfile/write/v4/ITsFileWriter.java | 4 ++++ .../apache/tsfile/write/v4/TsFileWriterBuilder.java | 4 ++++ .../org/apache/tsfile/write/writer/TsFileIOWriter.java | 2 ++ 22 files changed, 122 insertions(+), 3 deletions(-) diff --git a/java/common/src/main/java/org/apache/tsfile/enums/ColumnCategory.java b/java/common/src/main/java/org/apache/tsfile/enums/ColumnCategory.java index 3bfb8fd8e..d063cb213 100644 --- a/java/common/src/main/java/org/apache/tsfile/enums/ColumnCategory.java +++ b/java/common/src/main/java/org/apache/tsfile/enums/ColumnCategory.java @@ -22,6 +22,10 @@ import java.util.ArrayList; import java.util.List; +/** + * Column roles in table-model files. TAG columns identify entities, FIELD columns contain + * measurements, ATTRIBUTE columns contain descriptive values, and TIME is the timestamp column. + */ public enum ColumnCategory { TAG, FIELD, diff --git a/java/tsfile/src/main/java/org/apache/tsfile/encoding/encoder/SDTEncoder.java b/java/tsfile/src/main/java/org/apache/tsfile/encoding/encoder/SDTEncoder.java index 0915d12f0..f0ce4aecb 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/encoding/encoder/SDTEncoder.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/encoding/encoder/SDTEncoder.java @@ -19,6 +19,11 @@ package org.apache.tsfile.encoding.encoder; +/** + * Implements Swinging Door Trending lossy compression. compDeviation bounds value error, while + * compMinTime and compMaxTime control the time window. The pending boundary point is emitted when + * the door closes or flush is called. + */ public class SDTEncoder { // the last read time and value if upperDoor >= lowerDoor meaning out of compDeviation range, will diff --git a/java/tsfile/src/main/java/org/apache/tsfile/file/header/PageHeader.java b/java/tsfile/src/main/java/org/apache/tsfile/file/header/PageHeader.java index d752d509f..50e9be2d7 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/file/header/PageHeader.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/file/header/PageHeader.java @@ -53,6 +53,10 @@ public static int estimateMaxPageHeaderSizeWithoutStatistics() { return 2 * (Integer.BYTES + 1); } + /** + * Deserializes a page header from the input. This overload reads statistics only when + * hasStatistic is true; an empty page may have null statistics. + */ public static PageHeader deserializeFrom( InputStream inputStream, TSDataType dataType, boolean hasStatistic) throws IOException { int uncompressedSize = ReadWriteForEncodingUtils.readUnsignedVarInt(inputStream); @@ -173,7 +177,10 @@ public void setModified(boolean modified) { this.modified |= modified; } - /** max page header size without statistics. */ + /** + * Returns the serialized size of this page header and its page body, including statistics when + * statistics are present. + */ public int getSerializedPageSize() { if (uncompressedSize == 0) { // Empty page return ReadWriteForEncodingUtils.uVarIntSize(uncompressedSize); diff --git a/java/tsfile/src/main/java/org/apache/tsfile/file/metadata/MetadataIndexNode.java b/java/tsfile/src/main/java/org/apache/tsfile/file/metadata/MetadataIndexNode.java index 693b4ffb1..a24eef7e9 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/file/metadata/MetadataIndexNode.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/file/metadata/MetadataIndexNode.java @@ -125,6 +125,14 @@ public static MetadataIndexNode deserializeFrom( return new MetadataIndexNode(children, offset, nodeType); } + /** + * Finds the child entry for key. Exact search requires an equal key; non-exact search returns the + * floor entry. The selected entry's end offset is the exclusive upper bound of the child scan + * range. + * + * @param key -the sorted search key + * @param exactSearch -whether an exact key match is required + */ public Pair getChildIndexEntry(Comparable key, boolean exactSearch) { int index = binarySearchInChildren(key, exactSearch); if (index == -1) { diff --git a/java/tsfile/src/main/java/org/apache/tsfile/file/metadata/TimeseriesMetadata.java b/java/tsfile/src/main/java/org/apache/tsfile/file/metadata/TimeseriesMetadata.java index 790f457d1..ad905bf70 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/file/metadata/TimeseriesMetadata.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/file/metadata/TimeseriesMetadata.java @@ -49,6 +49,10 @@ import static org.apache.tsfile.utils.RamUsageEstimator.sizeOfCharArray; import static org.apache.tsfile.utils.RamUsageEstimator.sizeOfObjectArray; +/** + * Metadata for one time series. Chunk metadata may be loaded lazily from an in-memory buffer or + * temporary file; callers must initialize the loader before requesting deferred chunk metadata. + */ public class TimeseriesMetadata implements ITimeSeriesMetadata { private static final int INSTANCE_SIZE = diff --git a/java/tsfile/src/main/java/org/apache/tsfile/file/metadata/TsFileMetadata.java b/java/tsfile/src/main/java/org/apache/tsfile/file/metadata/TsFileMetadata.java index 95759ae20..3d8b683d6 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/file/metadata/TsFileMetadata.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/file/metadata/TsFileMetadata.java @@ -73,6 +73,9 @@ public static TsFileMetadata deserializeWithoutCacheTableSchemaMap( * deserialize data from the buffer. * * @param buffer -buffer use to deserialize + * @param context -reader context used for version and encryption compatibility + * @param needTableSchemaMap -whether table schemas should be materialized and cached; disabling + * it reduces memory usage. * @return -an instance of TsFileMetaData */ public static TsFileMetadata deserializeFrom( @@ -293,6 +296,10 @@ public Map getTableMetadataIndexNodeMap() { return tableMetadataIndexNodeMap; } + /** + * Returns the metadata index for tableName. If no table-specific node exists, the default root + * keyed by the empty table name is used. + */ public MetadataIndexNode getTableMetadataIndexNode(String tableName) { MetadataIndexNode metadataIndexNode = tableMetadataIndexNodeMap.get(tableName); if (metadataIndexNode == null) { diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/TsFileSequenceReader.java b/java/tsfile/src/main/java/org/apache/tsfile/read/TsFileSequenceReader.java index 2eea4089f..9f50888da 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/TsFileSequenceReader.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/TsFileSequenceReader.java @@ -160,6 +160,13 @@ public TsFileSequenceReader(String file) throws IOException { this(file, true, null); } + /** + * Scans the file and counts chunks in each chunk group. The current file position is used during + * scanning and the returned result may be incomplete when a malformed tail is encountered. + * + * @return the number of chunks for each scanned chunk group + * @throws IOException if the file cannot be read + */ public Map countChunksPerChunkGroup() throws IOException { Map result = new LinkedHashMap<>(); diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/common/Chunk.java b/java/tsfile/src/main/java/org/apache/tsfile/read/common/Chunk.java index d41ad14a3..c01884e73 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/common/Chunk.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/common/Chunk.java @@ -123,6 +123,12 @@ public void setDeleteIntervalList(List list) { this.deleteIntervalList = list; } + /** + * Appends the pages of another compatible chunk. When two single-page chunks are merged, the + * resulting chunk must be marked as multi-page and the original chunk statistics must be written + * into their page headers. Chunks must have compatible data types, codecs, and measurement + * semantics. + */ public void mergeChunkByAppendPage(Chunk chunk) throws IOException { int dataSize = 0; // from where the page data of the merged chunk starts, if -1, it means the merged chunk has @@ -216,6 +222,10 @@ public long getRetainedSizeInBytes() { return INSTANCE_SIZE + sizeOfByteArray(chunkData.capacity()); } + /** + * Re-encodes this chunk using the target data type and schema. Null values are represented using + * the target type's null sentinel, and statistics are recomputed from the rewritten records. + */ public Chunk rewrite(TSDataType newType, Chunk timeChunk) throws IOException { if (newType == null || newType == chunkHeader.getDataType()) { return this; diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/reader/IChunkReader.java b/java/tsfile/src/main/java/org/apache/tsfile/read/reader/IChunkReader.java index e13d32d97..93fd996bd 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/reader/IChunkReader.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/reader/IChunkReader.java @@ -24,6 +24,10 @@ import java.io.IOException; import java.util.List; +/** + * Provides page readers for one chunk and exposes chunk-level skip and deletion decisions. + * Implementations may defer decompression and decryption until a page is read. + */ public interface IChunkReader { boolean hasNextSatisfiedPage() throws IOException; diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/reader/IPageReader.java b/java/tsfile/src/main/java/org/apache/tsfile/read/reader/IPageReader.java index 68ce7ef32..5eeeecc63 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/reader/IPageReader.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/reader/IPageReader.java @@ -30,6 +30,10 @@ import java.util.List; import java.util.function.LongConsumer; +/** + * Reads decoded records from one page. Implementations may consume decoder state; filtering, + * deletion handling, pagination, and empty-page behavior are defined by the methods below. + */ public interface IPageReader extends IMetadata { default BatchData getAllSatisfiedPageData() throws IOException { diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/reader/chunk/ChunkReader.java b/java/tsfile/src/main/java/org/apache/tsfile/read/reader/chunk/ChunkReader.java index b555a25e1..cc32d2cbf 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/reader/chunk/ChunkReader.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/reader/chunk/ChunkReader.java @@ -79,6 +79,10 @@ public ChunkReader(Chunk chunk, long readStopTime) { this(chunk, readStopTime, null, null); } + /** + * Parses all page headers and creates page readers. For a single-page chunk, the chunk statistic + * may be reused; for multiple pages, each page follows its own statistic and data boundaries. + */ private void initAllPageReaders(Statistics chunkStatistic) { // construct next satisfied page header while (chunkDataBuffer.remaining() > 0) { @@ -105,6 +109,7 @@ private void initAllPageReaders(Statistics chunkStatisti } } + /** Determines whether a page can be skipped by time/statistic filters or deletion intervals. */ private boolean pageCanSkip(PageHeader pageHeader) { if (queryFilter != null && !queryFilter.satisfyStartEndTime(pageHeader.getStartTime(), pageHeader.getEndTime())) { @@ -116,6 +121,10 @@ private boolean pageCanSkip(PageHeader pageHeader) { return false; } + /** + * A fully deleted page is skipped; a partially deleted page remains readable and is marked for + * record-level filtering. + */ protected boolean pageDeleted(PageHeader pageHeader) { if (readStopTime > pageHeader.getEndTime()) { // used for chunk reader by timestamp diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/reader/page/PageReader.java b/java/tsfile/src/main/java/org/apache/tsfile/read/reader/page/PageReader.java index cdd45c4b2..f05bb09cc 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/reader/page/PageReader.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/reader/page/PageReader.java @@ -282,6 +282,10 @@ public void initTsBlockBuilder(List dataTypes) { // do nothing } + /** + * Tests whether timestamp is covered by the current deletion intervals. Intervals are inclusive + * and expected to be sorted; the deletion cursor advances monotonically for ascending timestamps. + */ protected boolean isDeleted(long timestamp) { while (deleteIntervalList != null && deleteCursor < deleteIntervalList.size()) { if (deleteIntervalList.get(deleteCursor).contains(timestamp)) { diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/v4/DeviceTableModelReader.java b/java/tsfile/src/main/java/org/apache/tsfile/read/v4/DeviceTableModelReader.java index a3d125a1d..68fde0744 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/v4/DeviceTableModelReader.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/v4/DeviceTableModelReader.java @@ -48,6 +48,11 @@ import java.util.Map; import java.util.Optional; +/** + * Reads table-model data for one device. Table and column names are normalized according to the + * reader's case rules; tagFilter applies only to tag columns. close() releases reader resources and + * may suppress close-time I/O failures. + */ public class DeviceTableModelReader implements ITsFileReader { protected TsFileSequenceReader fileReader; diff --git a/java/tsfile/src/main/java/org/apache/tsfile/write/TsFileWriter.java b/java/tsfile/src/main/java/org/apache/tsfile/write/TsFileWriter.java index 5521a4f72..653d0e2e4 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/write/TsFileWriter.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/write/TsFileWriter.java @@ -502,7 +502,11 @@ private void checkIsAllMeasurementsInGroup( } } - /** Check whether all measurements of dataPoints list are in the measurementGroup. */ + /** + * Validates measurement schemas against the registered group schema. For non-aligned writes, + * unknown measurements are filtered from a private copy; for aligned writes, an unknown + * measurement is rejected. The caller's list is not modified. + */ private List checkIsAllMeasurementsInGroup( List dataPoints, MeasurementGroup measurementGroup, boolean isAligned) throws NoMeasurementException { diff --git a/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/AlignedChunkGroupWriterImpl.java b/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/AlignedChunkGroupWriterImpl.java index 8896b58a4..be8a2bb1d 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/AlignedChunkGroupWriterImpl.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/AlignedChunkGroupWriterImpl.java @@ -50,6 +50,11 @@ import java.util.Set; import java.util.stream.Collectors; +/** + * Writes an aligned chunk group consisting of one shared time chunk and one value chunk per + * measurement. Time and value pages remain synchronized, and bitmaps represent missing values at + * aligned row positions. + */ public class AlignedChunkGroupWriterImpl implements IChunkGroupWriter { protected static final Logger LOG = LoggerFactory.getLogger(AlignedChunkGroupWriterImpl.class); @@ -278,6 +283,10 @@ public long getCurrentChunkGroupSize() { return size; } + /** + * Adds empty value pages and null bitmap entries when a value measurement is introduced after + * time pages already exist, preserving page boundaries and row alignment across all value chunks. + */ public void tryToAddEmptyPageAndData(ValueChunkWriter valueChunkWriter) throws IOException { // add empty page for (int i = 0; i < timeChunkWriter.getNumOfPages(); i++) { diff --git a/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/ChunkWriterImpl.java b/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/ChunkWriterImpl.java index 7d7912b14..b5a34fa52 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/ChunkWriterImpl.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/ChunkWriterImpl.java @@ -46,6 +46,11 @@ import java.nio.channels.WritableByteChannel; import java.util.function.Function; +/** + * Writes one measurement chunk by aggregating page data, serializing a chunk header and all page + * bodies, and producing chunk metadata. After the chunk is written, the writer resets its page + * state and can be reused. + */ public class ChunkWriterImpl implements IChunkWriter { private static final Logger logger = LoggerFactory.getLogger(ChunkWriterImpl.class); diff --git a/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/IChunkGroupWriter.java b/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/IChunkGroupWriter.java index de42c5585..bd24cac5c 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/IChunkGroupWriter.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/write/chunk/IChunkGroupWriter.java @@ -52,6 +52,12 @@ public interface IChunkGroupWriter { */ int write(Tablet tablet) throws WriteProcessException, IOException; + /** + * Writes rows in the half-open range [startRowIndex, endRowIndex). + * + * @param startRowIndex inclusive first row + * @param endRowIndex exclusive end row + */ int write(Tablet table, int startRowIndex, int endRowIndex) throws WriteProcessException, IOException; diff --git a/java/tsfile/src/main/java/org/apache/tsfile/write/schema/IMeasurementSchema.java b/java/tsfile/src/main/java/org/apache/tsfile/write/schema/IMeasurementSchema.java index ef83b4bbb..11dbb5d46 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/write/schema/IMeasurementSchema.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/write/schema/IMeasurementSchema.java @@ -32,6 +32,10 @@ import java.util.Map; import java.util.stream.Collectors; +/** + * Describes one measurement's name, data type, encoding, compression, statistic, and optional + * properties. Implementations must preserve these fields during schema serialization and copying. + */ public interface IMeasurementSchema extends Accountable { MeasurementSchemaType getSchemaType(); diff --git a/java/tsfile/src/main/java/org/apache/tsfile/write/schema/MeasurementSchema.java b/java/tsfile/src/main/java/org/apache/tsfile/write/schema/MeasurementSchema.java index 16dab7789..2f817477c 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/write/schema/MeasurementSchema.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/write/schema/MeasurementSchema.java @@ -72,7 +72,10 @@ public MeasurementSchema(String measurementName, TSDataType dataType) { null); } - /** set properties as an empty Map. */ + /** + * Creates a measurement schema. The properties argument may be null; callers must not assume that + * an empty mutable map is created automatically. + */ public MeasurementSchema(String measurementName, TSDataType dataType, TSEncoding encoding) { this( measurementName, diff --git a/java/tsfile/src/main/java/org/apache/tsfile/write/v4/ITsFileWriter.java b/java/tsfile/src/main/java/org/apache/tsfile/write/v4/ITsFileWriter.java index 809d1e5cf..f9a133d9b 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/write/v4/ITsFileWriter.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/write/v4/ITsFileWriter.java @@ -26,6 +26,10 @@ import java.io.IOException; +/** + * Public v4 TsFile writing contract. Implementations must document schema registration, + * record/tablet lifecycle, flush behavior, close requirements, and all checked exceptions. + */ public interface ITsFileWriter extends AutoCloseable { @TsFileApi diff --git a/java/tsfile/src/main/java/org/apache/tsfile/write/v4/TsFileWriterBuilder.java b/java/tsfile/src/main/java/org/apache/tsfile/write/v4/TsFileWriterBuilder.java index e5932eedd..dfbaacf63 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/write/v4/TsFileWriterBuilder.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/write/v4/TsFileWriterBuilder.java @@ -28,6 +28,10 @@ import java.io.File; import java.io.IOException; +/** + * Builder for v4 TsFile writers. build() validates the output path, schema, table-model settings, + * and memory thresholds before creating a writer. + */ public class TsFileWriterBuilder { private static final long defaultMemoryThresholdInByte = 32 * 1024 * 1024; diff --git a/java/tsfile/src/main/java/org/apache/tsfile/write/writer/TsFileIOWriter.java b/java/tsfile/src/main/java/org/apache/tsfile/write/writer/TsFileIOWriter.java index 93ddec569..dc2ee9d47 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/write/writer/TsFileIOWriter.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/write/writer/TsFileIOWriter.java @@ -331,8 +331,10 @@ public boolean isWritingChunkGroup() { * @param measurementId - measurementId of this time series * @param compressionCodecName - compression name of this time series * @param tsDataType - data type + * @param encodingType - the encoding used by the chunk pages * @param statistics - Chunk statistics * @param dataSize - the serialized size of all pages + * @param numOfPages - the number of serialized pages * @param mask - 0x80 for time chunk, 0x40 for value chunk, 0x00 for common chunk * @throws IOException if I/O error occurs */