This is an automated email from the ASF dual-hosted git repository. jiangtian pushed a commit to branch fix_compression_by_type in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 03c9c18cb40944531f82006aa01f9ae70e9502d7 Author: Tian Jiang <[email protected]> AuthorDate: Wed Aug 20 11:25:39 2025 +0800 Fix that compression by type is not properly applied --- .../iotdb/session/it/IoTDBSessionSimpleIT.java | 180 ++++++++++++++++++++- .../persistence/schema/TemplateTable.java | 2 +- .../org/apache/iotdb/db/conf/IoTDBDescriptor.java | 25 +++ .../analyze/schema/AutoCreateSchemaExecutor.java | 13 +- .../plan/analyze/schema/TemplateSchemaFetcher.java | 2 +- .../db/queryengine/plan/parser/ASTVisitor.java | 8 +- .../fetcher/TableHeaderSchemaValidator.java | 4 +- .../relational/sql/ast/WrappedInsertStatement.java | 2 +- .../metrics/IoTDBInternalLocalReporter.java | 3 +- .../FirstBatchCompactionAlignedChunkWriter.java | 3 +- 10 files changed, 227 insertions(+), 15 deletions(-) diff --git a/integration-test/src/test/java/org/apache/iotdb/session/it/IoTDBSessionSimpleIT.java b/integration-test/src/test/java/org/apache/iotdb/session/it/IoTDBSessionSimpleIT.java index 55487cbc847..c69d1438d73 100644 --- a/integration-test/src/test/java/org/apache/iotdb/session/it/IoTDBSessionSimpleIT.java +++ b/integration-test/src/test/java/org/apache/iotdb/session/it/IoTDBSessionSimpleIT.java @@ -2222,6 +2222,7 @@ public class IoTDBSessionSimpleIT { @Test public void testAlterDefaultCompression() throws IoTDBConnectionException, StatementExecutionException { + // auto-create try (ISession session = EnvFactory.getEnv().getSessionConnection()) { List<TSDataType> types = Arrays.asList( @@ -2288,7 +2289,7 @@ public class IoTDBSessionSimpleIT { String.format("SET CONFIGURATION '%s_compressor'='GZIP'", configName)); } - String device2 = "root.test.d1"; + String device2 = "root.test.d2"; session.insertRecord(device2, 0, measurements, types, values); try (SessionDataSet dataSet = @@ -2301,5 +2302,182 @@ public class IoTDBSessionSimpleIT { } } } + + // manual create + try (ISession session = EnvFactory.getEnv().getSessionConnection()) { + List<TSDataType> types = + Arrays.asList( + TSDataType.BOOLEAN, + TSDataType.INT32, + TSDataType.DATE, + TSDataType.INT64, + TSDataType.TIMESTAMP, + TSDataType.FLOAT, + TSDataType.DOUBLE, + TSDataType.TEXT, + TSDataType.STRING, + TSDataType.BLOB); + List<String> measurements = + types.stream().map(dataType -> "__" + dataType.toString()).collect(Collectors.toList()); + + String device3 = "root.test.d3"; + for (int i = 0; i < types.size(); i++) { + session.executeNonQueryStatement( + String.format( + "CREATE TIMESERIES %s.%s WITH DATATYPE=%s", + device3, measurements.get(i), types.get(i))); + } + + try (SessionDataSet dataSet = + session.executeQueryStatement("SHOW TIMESERIES root.test.d3.**")) { + int compressionIndex = dataSet.getColumnNames().indexOf("Compression"); + while (dataSet.hasNext()) { + RowRecord rec = dataSet.next(); + Field compressionField = rec.getFields().get(compressionIndex); + assertEquals("GZIP", compressionField.getStringValue()); + } + } + + for (TSDataType type : types) { + String configName = null; + switch (type) { + case INT32: + case INT64: + case FLOAT: + case DOUBLE: + case TEXT: + case BOOLEAN: + configName = type.name().toLowerCase(); + break; + case STRING: + case BLOB: + configName = "text"; + break; + case DATE: + configName = "int32"; + break; + case TIMESTAMP: + configName = "int64"; + break; + } + session.executeNonQueryStatement( + String.format("SET CONFIGURATION '%s_compressor'='LZ4'", configName)); + } + + String device4 = "root.test.d4"; + for (int i = 0; i < types.size(); i++) { + session.executeNonQueryStatement( + String.format( + "CREATE TIMESERIES %s.%s WITH DATATYPE=%s", + device4, measurements.get(i), types.get(i))); + } + + try (SessionDataSet dataSet = + session.executeQueryStatement("SHOW TIMESERIES root.test.d4.**")) { + int compressionIndex = dataSet.getColumnNames().indexOf("Compression"); + while (dataSet.hasNext()) { + RowRecord rec = dataSet.next(); + Field compressionField = rec.getFields().get(compressionIndex); + assertEquals("LZ4", compressionField.getStringValue()); + } + } + } + + // template + try (ISession session = EnvFactory.getEnv().getSessionConnection()) { + List<TSDataType> types = + Arrays.asList( + TSDataType.BOOLEAN, + TSDataType.INT32, + TSDataType.DATE, + TSDataType.INT64, + TSDataType.TIMESTAMP, + TSDataType.FLOAT, + TSDataType.DOUBLE, + TSDataType.TEXT, + TSDataType.STRING, + TSDataType.BLOB); + List<String> measurements = + types.stream().map(dataType -> "__" + dataType.toString()).collect(Collectors.toList()); + List<Object> values = + Arrays.asList( + false, + 1, + LocalDate.of(1000, 1, 1), + 1L, + 1L, + 1.0f, + 1.0, + new Binary("1".getBytes(StandardCharsets.UTF_8)), + new Binary("1".getBytes(StandardCharsets.UTF_8)), + new Binary("1".getBytes(StandardCharsets.UTF_8))); + + String createTemplateSql = "CREATE DEVICE TEMPLATE t1 ("; + for (int i = 0; i < types.size(); i++) { + createTemplateSql += measurements.get(i) + " " + types.get(i).name(); + if (i != types.size() - 1) { + createTemplateSql += ","; + } + } + createTemplateSql += ")"; + session.executeNonQueryStatement(createTemplateSql); + + session.executeNonQueryStatement("SET DEVICE TEMPLATE t1 TO root.test.d5"); + String device5 = "root.test.d5"; + session.insertRecord(device5, 0, measurements, types, values); + + try (SessionDataSet dataSet = + session.executeQueryStatement("SHOW TIMESERIES root.test.d5.**")) { + int compressionIndex = dataSet.getColumnNames().indexOf("Compression"); + while (dataSet.hasNext()) { + RowRecord rec = dataSet.next(); + Field compressionField = rec.getFields().get(compressionIndex); + assertEquals("LZ4", compressionField.getStringValue()); + } + } + + for (TSDataType type : types) { + String configName = null; + switch (type) { + case INT32: + case INT64: + case FLOAT: + case DOUBLE: + case TEXT: + case BOOLEAN: + configName = type.name().toLowerCase(); + break; + case STRING: + case BLOB: + configName = "text"; + break; + case DATE: + configName = "int32"; + break; + case TIMESTAMP: + configName = "int64"; + break; + } + session.executeNonQueryStatement( + String.format("SET CONFIGURATION '%s_compressor'='GZIP'", configName)); + } + + createTemplateSql = createTemplateSql.replace("t1", "t2"); + session.executeNonQueryStatement(createTemplateSql); + session.executeNonQueryStatement("SET DEVICE TEMPLATE t2 TO root.test.d6"); + + String device6 = "root.test.d6"; + session.insertRecord(device6, 0, measurements, types, values); + + try (SessionDataSet dataSet = + session.executeQueryStatement("SHOW TIMESERIES root.test.d6.**")) { + int compressionIndex = dataSet.getColumnNames().indexOf("Compression"); + while (dataSet.hasNext()) { + RowRecord rec = dataSet.next(); + Field compressionField = rec.getFields().get(compressionIndex); + assertEquals("GZIP", compressionField.getStringValue()); + } + } + } } } diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/schema/TemplateTable.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/schema/TemplateTable.java index 7557eacb6e5..6bb06d99db8 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/schema/TemplateTable.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/schema/TemplateTable.java @@ -156,7 +156,7 @@ public class TemplateTable { dataTypeList.get(i), encodingList == null ? getDefaultEncoding(dataTypeList.get(i)) : encodingList.get(i), compressionTypeList == null - ? TSFileDescriptor.getInstance().getConfig().getCompressor() + ? TSFileDescriptor.getInstance().getConfig().getCompressor(dataTypeList.get(i)) : compressionTypeList.get(i)); } else { if (!measurementSchema.getType().equals(dataTypeList.get(i)) diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java index f9377892ef3..f89e76b1dc3 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java @@ -1821,6 +1821,31 @@ public class IoTDBDescriptor { TSFileDescriptor.getInstance() .getConfig() .setEncryptType(properties.getProperty("encrypt_type", "UNENCRYPTED")); + + String booleanCompressor = properties.getProperty("boolean_compressor"); + if (booleanCompressor != null) { + TSFileDescriptor.getInstance().getConfig().setBooleanCompression(booleanCompressor); + } + String int32Compressor = properties.getProperty("int32_compressor"); + if (int32Compressor != null) { + TSFileDescriptor.getInstance().getConfig().setInt32Compression(int32Compressor); + } + String int64Compressor = properties.getProperty("int64_compressor"); + if (int64Compressor != null) { + TSFileDescriptor.getInstance().getConfig().setInt64Compression(int64Compressor); + } + String floatCompressor = properties.getProperty("float_compressor"); + if (floatCompressor != null) { + TSFileDescriptor.getInstance().getConfig().setFloatCompression(floatCompressor); + } + String doubleCompressor = properties.getProperty("double_compressor"); + if (doubleCompressor != null) { + TSFileDescriptor.getInstance().getConfig().setDoubleCompression(doubleCompressor); + } + String textCompressor = properties.getProperty("text_compressor"); + if (textCompressor != null) { + TSFileDescriptor.getInstance().getConfig().setTextCompression(textCompressor); + } } // Mqtt related diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/AutoCreateSchemaExecutor.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/AutoCreateSchemaExecutor.java index dd3c925786a..9b8d1047ce0 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/AutoCreateSchemaExecutor.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/AutoCreateSchemaExecutor.java @@ -123,7 +123,7 @@ class AutoCreateSchemaExecutor { dataTypesOfMissingMeasurement.add(tsDataType); encodingsOfMissingMeasurement.add(getDefaultEncoding(tsDataType)); compressionTypesOfMissingMeasurement.add( - TSFileDescriptor.getInstance().getConfig().getCompressor()); + TSFileDescriptor.getInstance().getConfig().getCompressor(tsDataType)); } }); @@ -180,7 +180,9 @@ class AutoCreateSchemaExecutor { measurements[measurementIndex], tsDataTypes[measurementIndex], getDefaultEncoding(tsDataTypes[measurementIndex]), - TSFileDescriptor.getInstance().getConfig().getCompressor()); + TSFileDescriptor.getInstance() + .getConfig() + .getCompressor(tsDataTypes[measurementIndex])); } return v; }); @@ -345,7 +347,9 @@ class AutoCreateSchemaExecutor { ? getDefaultEncoding(tsDataTypes[measurementIndex]) : encodings[measurementIndex], compressionTypes == null - ? TSFileDescriptor.getInstance().getConfig().getCompressor() + ? TSFileDescriptor.getInstance() + .getConfig() + .getCompressor(tsDataTypes[measurementIndex]) : compressionTypes[measurementIndex]); } return v; @@ -389,7 +393,8 @@ class AutoCreateSchemaExecutor { && compressionTypesList.get(finalDeviceIndex1) != null) { compressionType = compressionTypesList.get(finalDeviceIndex1)[index]; } else { - compressionType = TSFileDescriptor.getInstance().getConfig().getCompressor(); + compressionType = + TSFileDescriptor.getInstance().getConfig().getCompressor(dataType); } templateExtendInfo.addMeasurement( measurement, dataType, encoding, compressionType); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/TemplateSchemaFetcher.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/TemplateSchemaFetcher.java index 3516182bbf1..2a61b546d5e 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/TemplateSchemaFetcher.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/TemplateSchemaFetcher.java @@ -155,7 +155,7 @@ class TemplateSchemaFetcher { measurements[j], dataType, getDefaultEncoding(dataType), - TSFileDescriptor.getInstance().getConfig().getCompressor()); + TSFileDescriptor.getInstance().getConfig().getCompressor(dataType)); } } } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/parser/ASTVisitor.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/parser/ASTVisitor.java index 706d14f052c..328cb6c5e18 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/parser/ASTVisitor.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/parser/ASTVisitor.java @@ -461,7 +461,9 @@ public class ASTVisitor extends IoTDBSqlParserBaseVisitor<Statement> { } createTimeSeriesStatement.setCompressor( - TSFileDescriptor.getInstance().getConfig().getCompressor()); + TSFileDescriptor.getInstance() + .getConfig() + .getCompressor(createTimeSeriesStatement.getDataType())); if (props != null && props.containsKey(IoTDBConstant.COLUMN_TIMESERIES_COMPRESSION.toLowerCase())) { String compressionString = @@ -523,7 +525,7 @@ public class ASTVisitor extends IoTDBSqlParserBaseVisitor<Statement> { createAlignedTimeSeriesStatement.addEncoding(encoding); } - CompressionType compressor = TSFileDescriptor.getInstance().getConfig().getCompressor(); + CompressionType compressor = TSFileDescriptor.getInstance().getConfig().getCompressor(dataType); if (props.containsKey(IoTDBConstant.COLUMN_TIMESERIES_COMPRESSOR.toLowerCase())) { String compressorString = props.get(IoTDBConstant.COLUMN_TIMESERIES_COMPRESSOR.toLowerCase()).toUpperCase(); @@ -3745,7 +3747,7 @@ public class ASTVisitor extends IoTDBSqlParserBaseVisitor<Statement> { encodings.add(encoding); } - CompressionType compressor = TSFileDescriptor.getInstance().getConfig().getCompressor(); + CompressionType compressor = TSFileDescriptor.getInstance().getConfig().getCompressor(dataType); if (props.containsKey(IoTDBConstant.COLUMN_TIMESERIES_COMPRESSOR.toLowerCase())) { String compressorString = props.get(IoTDBConstant.COLUMN_TIMESERIES_COMPRESSOR.toLowerCase()).toUpperCase(); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/metadata/fetcher/TableHeaderSchemaValidator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/metadata/fetcher/TableHeaderSchemaValidator.java index 286a9d2dd70..2d127c3aeb5 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/metadata/fetcher/TableHeaderSchemaValidator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/metadata/fetcher/TableHeaderSchemaValidator.java @@ -338,7 +338,7 @@ public class TableHeaderSchemaValidator { columnName, dataType, getDefaultEncoding(dataType), - TSFileDescriptor.getInstance().getConfig().getCompressor()) + TSFileDescriptor.getInstance().getConfig().getCompressor(dataType)) // Unknown appears only for tree view field when the type needs auto-detection // Skip encoding & compressors because view query does not need these : new FieldColumnSchema(columnName, dataType); @@ -415,7 +415,7 @@ public class TableHeaderSchemaValidator { inputColumn.getName(), dataType, getDefaultEncoding(dataType), - TSFileDescriptor.getInstance().getConfig().getCompressor())); + TSFileDescriptor.getInstance().getConfig().getCompressor(dataType))); break; case TIME: throw new SemanticException( diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/WrappedInsertStatement.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/WrappedInsertStatement.java index 1cb351ec9f5..e0bd9cc3243 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/WrappedInsertStatement.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/WrappedInsertStatement.java @@ -236,7 +236,7 @@ public abstract class WrappedInsertStatement extends WrappedStatement real.getName(), tsDataType, getDefaultEncoding(tsDataType), - TSFileDescriptor.getInstance().getConfig().getCompressor()); + TSFileDescriptor.getInstance().getConfig().getCompressor(tsDataType)); innerTreeStatement.setMeasurementSchema(measurementSchema, i); try { innerTreeStatement.selfCheckDataTypes(i); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/metrics/IoTDBInternalLocalReporter.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/metrics/IoTDBInternalLocalReporter.java index 083d5fcd06c..c35ca1a704e 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/metrics/IoTDBInternalLocalReporter.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/metrics/IoTDBInternalLocalReporter.java @@ -220,7 +220,8 @@ public class IoTDBInternalLocalReporter extends IoTDBInternalReporter { TSDataType type = inferType(entry.getValue()); types.add(type.ordinal()); encodings.add((int) getDefaultEncoding(type).serialize()); - compressors.add((int) TSFileDescriptor.getInstance().getConfig().getCompressor().serialize()); + compressors.add( + (int) TSFileDescriptor.getInstance().getConfig().getCompressor(type).serialize()); } request.setPaths(paths); request.setDataTypes(types); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/executor/batch/utils/FirstBatchCompactionAlignedChunkWriter.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/executor/batch/utils/FirstBatchCompactionAlignedChunkWriter.java index 8413abaf74e..4d074c3c3a6 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/executor/batch/utils/FirstBatchCompactionAlignedChunkWriter.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/compaction/execute/utils/executor/batch/utils/FirstBatchCompactionAlignedChunkWriter.java @@ -101,7 +101,8 @@ public class FirstBatchCompactionAlignedChunkWriter extends AlignedChunkWriterIm TSEncoding timeEncoding = TSEncoding.valueOf(TSFileDescriptor.getInstance().getConfig().getTimeEncoder()); TSDataType timeType = TSFileDescriptor.getInstance().getConfig().getTimeSeriesDataType(); - CompressionType timeCompression = TSFileDescriptor.getInstance().getConfig().getCompressor(); + CompressionType timeCompression = + TSFileDescriptor.getInstance().getConfig().getCompressor(timeType); timeChunkWriter = new FirstBatchCompactionTimeChunkWriter( "",
