This is an automated email from the ASF dual-hosted git repository. jt2594838 pushed a commit to branch add_interfaces_for_type in repository https://gitbox.apache.org/repos/asf/tsfile.git
commit bab44db47169199bac9f30c2539ad4260a55b74e Author: Tian Jiang <[email protected]> AuthorDate: Thu Jul 23 11:06:17 2026 +0800 add Column.writeTo --- .../org/apache/tsfile/block/column/Column.java | 11 ++++ .../read/common/block/column/BinaryColumn.java | 15 ++++++ .../read/common/block/column/BooleanColumn.java | 11 ++++ .../read/common/block/column/DictionaryColumn.java | 13 +++++ .../read/common/block/column/DoubleColumn.java | 11 ++++ .../read/common/block/column/FloatColumn.java | 11 ++++ .../tsfile/read/common/block/column/IntColumn.java | 11 ++++ .../read/common/block/column/LongColumn.java | 11 ++++ .../read/common/block/column/NullColumn.java | 11 ++++ .../block/column/RunLengthEncodedColumn.java | 13 +++++ .../read/common/block/column/TimeColumn.java | 11 ++++ .../org/apache/tsfile/read/common/ColumnTest.java | 58 ++++++++++++++++++++++ 12 files changed, 187 insertions(+) diff --git a/java/common/src/main/java/org/apache/tsfile/block/column/Column.java b/java/common/src/main/java/org/apache/tsfile/block/column/Column.java index 9b7c88edd..445396f53 100644 --- a/java/common/src/main/java/org/apache/tsfile/block/column/Column.java +++ b/java/common/src/main/java/org/apache/tsfile/block/column/Column.java @@ -25,6 +25,7 @@ import org.apache.tsfile.utils.TsPrimitiveType; import java.io.DataOutputStream; import java.io.IOException; +import java.nio.ByteBuffer; import java.util.Arrays; public interface Column { @@ -110,6 +111,16 @@ public interface Column { throw new UnsupportedOperationException(getClass().getName()); } + /** Serializes the value at {@code index} to {@code buffer}. */ + default void writeTo(int index, ByteBuffer buffer) { + throw new UnsupportedOperationException(getClass().getName()); + } + + /** Serializes the value at {@code index} to {@code stream}. */ + default void writeTo(int index, DataOutputStream stream) throws IOException { + throw new UnsupportedOperationException(getClass().getName()); + } + /** Serializes all non-null values in this column. */ default void serializeWithoutNulls(DataOutputStream output) throws IOException { throw new UnsupportedOperationException(getClass().getName()); diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/BinaryColumn.java b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/BinaryColumn.java index a11885f7a..07a2b58db 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/BinaryColumn.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/BinaryColumn.java @@ -29,6 +29,7 @@ import org.apache.tsfile.utils.TsPrimitiveType; import java.io.DataOutputStream; import java.io.IOException; +import java.nio.ByteBuffer; import java.util.Arrays; import java.util.Optional; @@ -145,6 +146,20 @@ public class BinaryColumn implements Column { return new TsPrimitiveType.TsBinary(getBinary(position)); } + @Override + public void writeTo(int index, ByteBuffer buffer) { + Binary value = values[index + arrayOffset]; + buffer.putInt(value.getLength()); + buffer.put(value.getValues()); + } + + @Override + public void writeTo(int index, DataOutputStream stream) throws IOException { + Binary value = values[index + arrayOffset]; + stream.writeInt(value.getLength()); + stream.write(value.getValues()); + } + @Override public void serializeWithoutNulls(DataOutputStream output) throws IOException { for (int i = 0; i < positionCount; i++) { diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/BooleanColumn.java b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/BooleanColumn.java index ac0c0dac1..7d30c6a00 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/BooleanColumn.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/BooleanColumn.java @@ -29,6 +29,7 @@ import org.apache.tsfile.utils.TsPrimitiveType; import java.io.DataOutputStream; import java.io.IOException; +import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; import java.util.Arrays; import java.util.Optional; @@ -127,6 +128,16 @@ public class BooleanColumn implements Column { return new TsPrimitiveType.TsBoolean(getBoolean(position)); } + @Override + public void writeTo(int index, ByteBuffer buffer) { + buffer.put(values[index + arrayOffset] ? (byte) 1 : (byte) 0); + } + + @Override + public void writeTo(int index, DataOutputStream stream) throws IOException { + stream.writeBoolean(values[index + arrayOffset]); + } + @Override public void serializeWithoutNulls(DataOutputStream output) throws IOException { int nonNullCount = 0; diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/DictionaryColumn.java b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/DictionaryColumn.java index e26f42c09..1e2f7067b 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/DictionaryColumn.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/DictionaryColumn.java @@ -27,6 +27,9 @@ import org.apache.tsfile.utils.Binary; import org.apache.tsfile.utils.RamUsageEstimator; import org.apache.tsfile.utils.TsPrimitiveType; +import java.io.DataOutputStream; +import java.io.IOException; +import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Arrays; import java.util.HashMap; @@ -427,6 +430,16 @@ public final class DictionaryColumn implements Column { return dictionary.getTsPrimitiveType(position); } + @Override + public void writeTo(int index, ByteBuffer buffer) { + dictionary.writeTo(getId(index), buffer); + } + + @Override + public void writeTo(int index, DataOutputStream stream) throws IOException { + dictionary.writeTo(getId(index), stream); + } + @Override public boolean arePositionsEqual(int thisPos, Column that, int thatPos) { return dictionary.arePositionsEqual(getId(thisPos), that, thatPos); diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/DoubleColumn.java b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/DoubleColumn.java index 366a99d42..81e19a307 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/DoubleColumn.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/DoubleColumn.java @@ -29,6 +29,7 @@ import org.apache.tsfile.utils.TsPrimitiveType; import java.io.DataOutputStream; import java.io.IOException; +import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; import java.util.Arrays; import java.util.Optional; @@ -128,6 +129,16 @@ public class DoubleColumn implements Column { return new TsPrimitiveType.TsDouble(getDouble(position)); } + @Override + public void writeTo(int index, ByteBuffer buffer) { + buffer.putDouble(values[index + arrayOffset]); + } + + @Override + public void writeTo(int index, DataOutputStream stream) throws IOException { + stream.writeDouble(values[index + arrayOffset]); + } + @Override public void serializeWithoutNulls(DataOutputStream output) throws IOException { for (int i = 0; i < positionCount; i++) { diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/FloatColumn.java b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/FloatColumn.java index 35b110cb8..56ce7a1f8 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/FloatColumn.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/FloatColumn.java @@ -29,6 +29,7 @@ import org.apache.tsfile.utils.TsPrimitiveType; import java.io.DataOutputStream; import java.io.IOException; +import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; import java.util.Arrays; import java.util.Optional; @@ -143,6 +144,16 @@ public class FloatColumn implements Column { return new TsPrimitiveType.TsFloat(getFloat(position)); } + @Override + public void writeTo(int index, ByteBuffer buffer) { + buffer.putFloat(values[index + arrayOffset]); + } + + @Override + public void writeTo(int index, DataOutputStream stream) throws IOException { + stream.writeFloat(values[index + arrayOffset]); + } + @Override public void serializeWithoutNulls(DataOutputStream output) throws IOException { for (int i = 0; i < positionCount; i++) { diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/IntColumn.java b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/IntColumn.java index 10e6179d3..59d3a7505 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/IntColumn.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/IntColumn.java @@ -29,6 +29,7 @@ import org.apache.tsfile.utils.TsPrimitiveType; import java.io.DataOutputStream; import java.io.IOException; +import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; import java.util.Arrays; import java.util.Optional; @@ -199,6 +200,16 @@ public class IntColumn implements Column { return new TsPrimitiveType.TsInt(getInt(position), dataType); } + @Override + public void writeTo(int index, ByteBuffer buffer) { + buffer.putInt(values[index + arrayOffset]); + } + + @Override + public void writeTo(int index, DataOutputStream stream) throws IOException { + stream.writeInt(values[index + arrayOffset]); + } + @Override public void serializeWithoutNulls(DataOutputStream output) throws IOException { for (int i = 0; i < positionCount; i++) { diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/LongColumn.java b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/LongColumn.java index 8f089cdf2..2b6bc46fc 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/LongColumn.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/LongColumn.java @@ -29,6 +29,7 @@ import org.apache.tsfile.utils.TsPrimitiveType; import java.io.DataOutputStream; import java.io.IOException; +import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; import java.util.Arrays; import java.util.Optional; @@ -143,6 +144,16 @@ public class LongColumn implements Column { return new TsPrimitiveType.TsLong(getLong(position)); } + @Override + public void writeTo(int index, ByteBuffer buffer) { + buffer.putLong(values[index + arrayOffset]); + } + + @Override + public void writeTo(int index, DataOutputStream stream) throws IOException { + stream.writeLong(values[index + arrayOffset]); + } + @Override public void serializeWithoutNulls(DataOutputStream output) throws IOException { for (int i = 0; i < positionCount; i++) { diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/NullColumn.java b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/NullColumn.java index 06d720275..0d0145295 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/NullColumn.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/NullColumn.java @@ -27,6 +27,7 @@ import org.apache.tsfile.read.common.type.Type; import org.apache.tsfile.utils.RamUsageEstimator; import java.io.DataOutputStream; +import java.nio.ByteBuffer; import static java.util.Objects.requireNonNull; import static org.apache.tsfile.read.common.block.column.ColumnUtil.checkArrayRange; @@ -89,6 +90,16 @@ public class NullColumn implements Column { // There are no non-null values to serialize. } + @Override + public void writeTo(int index, ByteBuffer buffer) { + // There is no value to serialize. + } + + @Override + public void writeTo(int index, DataOutputStream stream) { + // There is no value to serialize. + } + @Override public int getPositionCount() { return positionCount; diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/RunLengthEncodedColumn.java b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/RunLengthEncodedColumn.java index 548cfab81..30c3af73e 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/RunLengthEncodedColumn.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/RunLengthEncodedColumn.java @@ -27,6 +27,9 @@ import org.apache.tsfile.utils.Binary; import org.apache.tsfile.utils.RamUsageEstimator; import org.apache.tsfile.utils.TsPrimitiveType; +import java.io.DataOutputStream; +import java.io.IOException; +import java.nio.ByteBuffer; import java.util.Arrays; import static java.util.Objects.requireNonNull; @@ -165,6 +168,16 @@ public class RunLengthEncodedColumn implements Column { return value.getTsPrimitiveType(0); } + @Override + public void writeTo(int index, ByteBuffer buffer) { + value.writeTo(0, buffer); + } + + @Override + public void writeTo(int index, DataOutputStream stream) throws IOException { + value.writeTo(0, stream); + } + @Override public boolean arePositionsEqual(int thisPos, Column that, int thatPos) { return value.arePositionsEqual(0, that, thatPos); diff --git a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/TimeColumn.java b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/TimeColumn.java index 4156482af..c2a5e7427 100644 --- a/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/TimeColumn.java +++ b/java/tsfile/src/main/java/org/apache/tsfile/read/common/block/column/TimeColumn.java @@ -27,6 +27,7 @@ import org.apache.tsfile.utils.RamUsageEstimator; import java.io.DataOutputStream; import java.io.IOException; +import java.nio.ByteBuffer; import java.util.Arrays; import static org.apache.tsfile.read.common.block.column.ColumnUtil.checkArrayRange; @@ -92,6 +93,16 @@ public class TimeColumn implements Column { return getLong(position); } + @Override + public void writeTo(int index, ByteBuffer buffer) { + buffer.putLong(values[index + arrayOffset]); + } + + @Override + public void writeTo(int index, DataOutputStream stream) throws IOException { + stream.writeLong(values[index + arrayOffset]); + } + @Override public void serializeWithoutNulls(DataOutputStream output) throws IOException { for (int i = 0; i < positionCount; i++) { diff --git a/java/tsfile/src/test/java/org/apache/tsfile/read/common/ColumnTest.java b/java/tsfile/src/test/java/org/apache/tsfile/read/common/ColumnTest.java index 6d6fe91a1..5e71a89c3 100644 --- a/java/tsfile/src/test/java/org/apache/tsfile/read/common/ColumnTest.java +++ b/java/tsfile/src/test/java/org/apache/tsfile/read/common/ColumnTest.java @@ -38,17 +38,75 @@ import org.apache.tsfile.read.common.block.column.NullColumn; import org.apache.tsfile.read.common.block.column.RunLengthEncodedColumn; import org.apache.tsfile.read.common.block.column.TimeColumn; import org.apache.tsfile.read.common.block.column.TimeColumnBuilder; +import org.apache.tsfile.read.common.type.Type; import org.apache.tsfile.utils.Binary; import org.apache.tsfile.utils.BytesUtils; import org.junit.Assert; import org.junit.Test; +import java.io.ByteArrayOutputStream; +import java.io.DataOutputStream; +import java.io.IOException; +import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; +import java.util.Arrays; import java.util.Optional; public class ColumnTest { + @Test + public void testWriteTo() throws IOException { + Binary binary = new Binary("test", StandardCharsets.UTF_8); + Column intDictionary = + DictionaryColumn.create( + 2, new IntColumn(3, Optional.empty(), new int[] {7, 8, 9}), new int[] {2, 0}); + Column binaryRle = + new RunLengthEncodedColumn(new BinaryColumn(1, Optional.empty(), new Binary[] {binary}), 3); + + Object[][] testCases = { + {new BooleanColumn(2, Optional.empty(), new boolean[] {false, true}).getRegion(1, 1), 0}, + {new IntColumn(2, Optional.empty(), new int[] {0, 42}).getRegion(1, 1), 0}, + {new LongColumn(2, Optional.empty(), new long[] {0L, 42L}).getRegion(1, 1), 0}, + {new FloatColumn(2, Optional.empty(), new float[] {0.0F, 1.25F}).getRegion(1, 1), 0}, + {new DoubleColumn(2, Optional.empty(), new double[] {0.0D, 2.5D}).getRegion(1, 1), 0}, + { + new BinaryColumn(2, Optional.empty(), new Binary[] {Binary.EMPTY_VALUE, binary}) + .getRegion(1, 1), + 0 + }, + {new TimeColumn(2, new long[] {0L, 42L}).getRegion(1, 1), 0}, + {intDictionary, 0}, + {binaryRle, 2} + }; + + for (Object[] testCase : testCases) { + Column column = (Column) testCase[0]; + int index = (int) testCase[1]; + ByteArrayOutputStream expectedOutput = new ByteArrayOutputStream(); + Type.fromTsDataType(column.getDataType()) + .serializeValue(column.getObject(index), new DataOutputStream(expectedOutput)); + byte[] expected = expectedOutput.toByteArray(); + + ByteBuffer buffer = ByteBuffer.allocate(expected.length); + column.writeTo(index, buffer); + Assert.assertEquals(expected.length, buffer.position()); + Assert.assertArrayEquals(expected, Arrays.copyOf(buffer.array(), buffer.position())); + + ByteArrayOutputStream actualOutput = new ByteArrayOutputStream(); + column.writeTo(index, new DataOutputStream(actualOutput)); + Assert.assertArrayEquals(expected, actualOutput.toByteArray()); + } + + Column nullColumn = new NullColumn(1); + ByteBuffer emptyBuffer = ByteBuffer.allocate(0); + nullColumn.writeTo(0, emptyBuffer); + ByteArrayOutputStream emptyOutput = new ByteArrayOutputStream(); + nullColumn.writeTo(0, new DataOutputStream(emptyOutput)); + Assert.assertEquals(0, emptyBuffer.position()); + Assert.assertEquals(0, emptyOutput.size()); + } + @Test public void testArePositionsEqual() { IntColumnBuilder builder = new IntColumnBuilder(null, 5);
