This is an automated email from the ASF dual-hosted git repository.
jiangtian pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new e90e2326304 Fix that compression by type is not properly applied
(#16211)
e90e2326304 is described below
commit e90e2326304739be2d9692ae60795564de56227e
Author: Jiang Tian <[email protected]>
AuthorDate: Wed Aug 20 16:39:02 2025 +0800
Fix that compression by type is not properly applied (#16211)
---
.../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 17e01d1e340..440d08bbe84 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 aa8a8621784..26830571f4b 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
@@ -462,7 +462,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 =
@@ -524,7 +526,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();
@@ -3746,7 +3748,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(
"",