This is an automated email from the ASF dual-hosted git repository.

jackietien 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 4cc6e7f3971 Append time column at the end of input columns in mappable 
UDTF
4cc6e7f3971 is described below

commit 4cc6e7f3971bad9c1b07baaca23c84c43c41fdd9
Author: Zhihao Shen <[email protected]>
AuthorDate: Wed Mar 20 15:04:58 2024 +0800

    Append time column at the end of input columns in mappable UDTF
---
 .../org/apache/iotdb/udf/api/utils/RowImpl.java    |  3 ++-
 .../visitor/ColumnTransformerVisitor.java          | 20 +++++++++++++++---
 .../commons/udf/builtin/String/UDTFConcat.java     |  6 +++---
 .../iotdb/commons/udf/builtin/UDTFConst.java       | 24 +++++++++++-----------
 .../iotdb/commons/udf/builtin/UDTFConstE.java      |  8 ++++----
 .../iotdb/commons/udf/builtin/UDTFConstPi.java     |  8 ++++----
 .../read/common/block/column/TimeColumn.java       |  5 +++++
 7 files changed, 47 insertions(+), 27 deletions(-)

diff --git 
a/iotdb-api/udf-api/src/main/java/org/apache/iotdb/udf/api/utils/RowImpl.java 
b/iotdb-api/udf-api/src/main/java/org/apache/iotdb/udf/api/utils/RowImpl.java
index 94d6542b5f0..654081400c8 100644
--- 
a/iotdb-api/udf-api/src/main/java/org/apache/iotdb/udf/api/utils/RowImpl.java
+++ 
b/iotdb-api/udf-api/src/main/java/org/apache/iotdb/udf/api/utils/RowImpl.java
@@ -38,7 +38,8 @@ public class RowImpl implements Row {
 
   @Override
   public long getTime() {
-    return (long) rowRecord[size];
+    // Time column is always the last column
+    return (long) rowRecord[size - 1];
   }
 
   @Override
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/expression/visitor/ColumnTransformerVisitor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/expression/visitor/ColumnTransformerVisitor.java
index 75a419988a3..e0c7f18831b 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/expression/visitor/ColumnTransformerVisitor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/expression/visitor/ColumnTransformerVisitor.java
@@ -70,6 +70,7 @@ import 
org.apache.iotdb.db.queryengine.transformation.dag.udf.UDTFContext;
 import org.apache.iotdb.db.queryengine.transformation.dag.udf.UDTFExecutor;
 import org.apache.iotdb.db.queryengine.transformation.dag.util.TransformUtils;
 import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
+import org.apache.iotdb.tsfile.read.common.type.LongType;
 import org.apache.iotdb.tsfile.read.common.type.Type;
 import org.apache.iotdb.tsfile.read.common.type.TypeFactory;
 
@@ -235,9 +236,22 @@ public class ColumnTransformerVisitor
               functionExpression, 
getBuiltInScalarFunctionTransformer(functionExpression, context));
         } else {
           ColumnTransformer[] inputColumnTransformers =
-              expressions.stream()
-                  .map(expression -> this.process(expression, context))
-                  .toArray(ColumnTransformer[]::new);
+              new ColumnTransformer[expressions.size() + 1];
+          for (int i = 0; i < expressions.size(); i++) {
+            inputColumnTransformers[i] = this.process(expressions.get(i), 
context);
+          }
+          // Append time column at the end of input columns for mappable UDTF
+          ColumnTransformer columnTransformer =
+              context.cache.computeIfAbsent(
+                  new TimestampOperand(),
+                  e -> {
+                    TimeColumnTransformer timeColumnTransformer =
+                        new TimeColumnTransformer(LongType.getInstance());
+                    context.leafList.add(timeColumnTransformer);
+                    return timeColumnTransformer;
+                  });
+          columnTransformer.addReferenceCount();
+          inputColumnTransformers[expressions.size()] = columnTransformer;
 
           UDTFExecutor executor =
               
context.udtfContext.getExecutorByFunctionExpression(functionExpression);
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/String/UDTFConcat.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/String/UDTFConcat.java
index c63787f6e59..b84305c1f71 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/String/UDTFConcat.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/String/UDTFConcat.java
@@ -104,18 +104,18 @@ public class UDTFConcat implements UDTF {
     int rowCount = columns[0].getPositionCount();
 
     Binary[][] inputFrame = new Binary[colCount][rowCount];
-    for (int i = 0; i < colCount; i++) {
+    for (int i = 0; i < colCount - 1; i++) {
       inputFrame[i] = columns[i].getBinaries();
     }
 
     boolean[][] isNullFrame = new boolean[colCount][rowCount];
-    for (int i = 0; i < colCount; i++) {
+    for (int i = 0; i < colCount - 1; i++) {
       isNullFrame[i] = columns[i].isNull();
     }
 
     for (int row = 0; row < rowCount; row++) {
       StringBuilder concatSeries = new StringBuilder();
-      for (int col = 0; col < colCount; col++) {
+      for (int col = 0; col < colCount - 1; col++) {
         if (isNullFrame[col][row]) {
           continue;
         }
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFConst.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFConst.java
index 9c335262ad7..855433ec376 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFConst.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFConst.java
@@ -157,8 +157,8 @@ public class UDTFConst implements UDTF {
       case INT32:
         for (int i = 0; i < count; i++) {
           boolean hasWritten = false;
-          for (Column column : columns) {
-            if (!column.isNull(i)) {
+          for (int j = 0; j < columns.length - 1; j++) {
+            if (!columns[j].isNull(i)) {
               builder.writeInt(intValue);
               hasWritten = true;
               break;
@@ -172,8 +172,8 @@ public class UDTFConst implements UDTF {
       case INT64:
         for (int i = 0; i < count; i++) {
           boolean hasWritten = false;
-          for (Column column : columns) {
-            if (!column.isNull(i)) {
+          for (int j = 0; j < columns.length - 1; j++) {
+            if (!columns[j].isNull(i)) {
               builder.writeLong(longValue);
               hasWritten = true;
               break;
@@ -187,8 +187,8 @@ public class UDTFConst implements UDTF {
       case FLOAT:
         for (int i = 0; i < count; i++) {
           boolean hasWritten = false;
-          for (Column column : columns) {
-            if (!column.isNull(i)) {
+          for (int j = 0; j < columns.length - 1; j++) {
+            if (!columns[j].isNull(i)) {
               builder.writeFloat(floatValue);
               hasWritten = true;
               break;
@@ -202,8 +202,8 @@ public class UDTFConst implements UDTF {
       case DOUBLE:
         for (int i = 0; i < count; i++) {
           boolean hasWritten = false;
-          for (Column column : columns) {
-            if (!column.isNull(i)) {
+          for (int j = 0; j < columns.length - 1; j++) {
+            if (!columns[j].isNull(i)) {
               builder.writeDouble(doubleValue);
               hasWritten = true;
               break;
@@ -217,8 +217,8 @@ public class UDTFConst implements UDTF {
       case BOOLEAN:
         for (int i = 0; i < count; i++) {
           boolean hasWritten = false;
-          for (Column column : columns) {
-            if (!column.isNull(i)) {
+          for (int j = 0; j < columns.length - 1; j++) {
+            if (!columns[j].isNull(i)) {
               builder.writeBoolean(booleanValue);
               hasWritten = true;
               break;
@@ -232,8 +232,8 @@ public class UDTFConst implements UDTF {
       case TEXT:
         for (int i = 0; i < count; i++) {
           boolean hasWritten = false;
-          for (Column column : columns) {
-            if (!column.isNull(i)) {
+          for (int j = 0; j < columns.length - 1; j++) {
+            if (!columns[j].isNull(i)) {
               builder.writeBinary(binaryValue);
               hasWritten = true;
               break;
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFConstE.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFConstE.java
index 7be88d38a72..3e25f77a7cb 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFConstE.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFConstE.java
@@ -52,12 +52,12 @@ public class UDTFConstE implements UDTF {
 
   @Override
   public void transform(Column[] columns, ColumnBuilder builder) throws 
Exception {
-    int colCount = columns[0].getPositionCount();
+    int rowCount = columns[0].getPositionCount();
 
-    for (int i = 0; i < colCount; i++) {
+    for (int i = 0; i < rowCount; i++) {
       boolean hasWritten = false;
-      for (Column column : columns) {
-        if (!column.isNull(i)) {
+      for (int j = 0; j < columns.length - 1; j++) {
+        if (!columns[j].isNull(i)) {
           builder.writeDouble(Math.E);
           hasWritten = true;
           break;
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFConstPi.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFConstPi.java
index 51b0c7391a5..c25512c554a 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFConstPi.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/udf/builtin/UDTFConstPi.java
@@ -52,12 +52,12 @@ public class UDTFConstPi implements UDTF {
 
   @Override
   public void transform(Column[] columns, ColumnBuilder builder) throws 
Exception {
-    int colCount = columns[0].getPositionCount();
+    int rowCount = columns[0].getPositionCount();
 
-    for (int i = 0; i < colCount; i++) {
+    for (int i = 0; i < rowCount; i++) {
       boolean hasWritten = false;
-      for (Column column : columns) {
-        if (!column.isNull(i)) {
+      for (int j = 0; j < columns.length - 1; j++) {
+        if (!columns[j].isNull(i)) {
           builder.writeDouble(Math.PI);
           hasWritten = true;
           break;
diff --git 
a/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/TimeColumn.java
 
b/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/TimeColumn.java
index a18d5fcf3dd..c1c50208357 100644
--- 
a/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/TimeColumn.java
+++ 
b/iotdb-core/tsfile/src/main/java/org/apache/iotdb/tsfile/read/common/block/column/TimeColumn.java
@@ -144,6 +144,11 @@ public class TimeColumn implements Column {
     return values;
   }
 
+  @Override
+  public long[] getLongs() {
+    return getTimes();
+  }
+
   @Override
   public int getInstanceSize() {
     return INSTANCE_SIZE;

Reply via email to