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(
             "",

Reply via email to