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;