Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,12 @@ public Object deserialize(JsonParser parser, DeserializationContext context)
}
}

return unsupported(node, context);
}

/** Inputs a subclass accepts beyond strings and {@link FieldRef}s. */
protected Object unsupported(JsonNode node, DeserializationContext context)
throws java.io.IOException {
context.reportInputMismatch(
Object.class, "Unsupported StringTransform input JSON: %s", node.toString());
return null;
Expand All @@ -108,15 +114,7 @@ public final List<Object> inputs() {

@JsonGetter(FIELD_INPUTS)
public final List<Object> inputsForJson() {
List<Object> serialized = new ArrayList<>(inputs.size());
for (Object input : inputs) {
if (input instanceof BinaryString) {
serialized.add(input.toString());
} else {
serialized.add(input);
}
}
return serialized;
return inputsForJson(inputs);
}

@Override
Expand Down Expand Up @@ -157,8 +155,27 @@ public int hashCode() {

@Override
public String toString() {
List<String> inputs =
this.inputs.stream().map(Object::toString).collect(Collectors.toList());
return name() + "(" + String.join(", ", inputs) + ')';
return formatCall(name(), inputs);
}

/** Inputs as written to JSON: {@link BinaryString} literals become JSON strings. */
static List<Object> inputsForJson(List<Object> inputs) {
List<Object> serialized = new ArrayList<>(inputs.size());
for (Object input : inputs) {
if (input instanceof BinaryString) {
serialized.add(input.toString());
} else {
serialized.add(input);
}
}
return serialized;
}

/** Renders a transform as {@code NAME(input, input)}. */
static String formatCall(String name, List<Object> inputs) {
return name
+ "("
+ inputs.stream().map(String::valueOf).collect(Collectors.joining(", "))
+ ')';
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,15 @@
import org.apache.paimon.types.DataType;
import org.apache.paimon.types.DataTypes;

import org.apache.paimon.shade.jackson2.com.fasterxml.jackson.annotation.JsonCreator;
import org.apache.paimon.shade.jackson2.com.fasterxml.jackson.annotation.JsonGetter;
import org.apache.paimon.shade.jackson2.com.fasterxml.jackson.annotation.JsonIgnore;
import org.apache.paimon.shade.jackson2.com.fasterxml.jackson.annotation.JsonProperty;
import org.apache.paimon.shade.jackson2.com.fasterxml.jackson.databind.DeserializationContext;
import org.apache.paimon.shade.jackson2.com.fasterxml.jackson.databind.JsonNode;
import org.apache.paimon.shade.jackson2.com.fasterxml.jackson.databind.annotation.JsonDeserialize;

import java.io.IOException;
import java.util.List;
import java.util.Objects;

Expand All @@ -39,11 +48,37 @@ public class SubstringTransform implements Transform {

private final List<Object> inputs;

public SubstringTransform(List<Object> inputs) {
@JsonCreator
public SubstringTransform(
@JsonProperty(StringTransform.FIELD_INPUTS)
@JsonDeserialize(contentUsing = InputDeserializer.class)
List<Object> inputs) {
checkArgument(inputs.size() == 2 || inputs.size() == 3);
this.inputs = inputs;
}

/** Deserializer for {@link SubstringTransform} inputs, which may also be integers. */
public static class InputDeserializer extends StringTransform.InputDeserializer {

private static final long serialVersionUID = 1L;

@Override
protected Object unsupported(JsonNode node, DeserializationContext context)
throws IOException {
if (node.isNumber()) {
// canConvertToInt checks the range but not integrality
if (!node.isIntegralNumber()) {
context.reportInputMismatch(
Object.class,
"SubstringTransform position must be an integer: %s",
node.toString());
}
return node.canConvertToInt() ? node.intValue() : node.numberValue();
}
return super.unsupported(node, context);
}
}

@Override
public String name() {
return NAME;
Expand Down Expand Up @@ -71,6 +106,10 @@ public final Object transform(InternalRow row) {
if (begin instanceof FieldRef) {
FieldRef beginRef = (FieldRef) begin;
checkArgument(beginRef.type().is(INTEGER_NUMERIC));
// getInt on a null reads an undefined value on columnar rows
if (row.isNullAt(beginRef.index())) {
return null;
}
beginIndex = row.getInt(beginRef.index());
} else {
beginIndex = Integer.parseInt(inputs.get(1).toString());
Expand All @@ -85,6 +124,9 @@ public final Object transform(InternalRow row) {
if (end instanceof FieldRef) {
FieldRef endRef = (FieldRef) inputs.get(2);
checkArgument(endRef.type().is(INTEGER_NUMERIC));
if (row.isNullAt(endRef.index())) {
return null;
}
endIndex = beginIndex + row.getInt(endRef.index()) - 1;
} else {
endIndex = beginIndex + Integer.parseInt(inputs.get(2).toString()) - 1;
Expand All @@ -103,10 +145,16 @@ public Transform copyWithNewInputs(List<Object> inputs) {
}

@Override
@JsonIgnore
public final List<Object> inputs() {
return inputs;
}

@JsonGetter(StringTransform.FIELD_INPUTS)
public final List<Object> inputsForJson() {
return StringTransform.inputsForJson(inputs);
}

@Override
public boolean equals(Object o) {
if (o == null || getClass() != o.getClass()) {
Expand All @@ -125,4 +173,9 @@ public DataType outputType() {
public int hashCode() {
return Objects.hashCode(inputs);
}

@Override
public String toString() {
return StringTransform.formatCall(name(), inputs);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,8 @@
@JsonSubTypes.Type(value = ConcatWsTransform.class, name = ConcatWsTransform.NAME),
@JsonSubTypes.Type(value = UpperTransform.class, name = UpperTransform.NAME),
@JsonSubTypes.Type(value = LowerTransform.class, name = LowerTransform.NAME),
@JsonSubTypes.Type(value = SubstringTransform.class, name = SubstringTransform.NAME),
@JsonSubTypes.Type(value = TrimTransform.class, name = TrimTransform.NAME),
@JsonSubTypes.Type(value = NullTransform.class, name = NullTransform.NAME)
})
public interface Transform extends Serializable {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,9 +21,16 @@
import org.apache.paimon.data.BinaryString;
import org.apache.paimon.utils.StringUtils;

import org.apache.paimon.shade.jackson2.com.fasterxml.jackson.annotation.JsonCreator;
import org.apache.paimon.shade.jackson2.com.fasterxml.jackson.annotation.JsonGetter;
import org.apache.paimon.shade.jackson2.com.fasterxml.jackson.annotation.JsonProperty;
import org.apache.paimon.shade.jackson2.com.fasterxml.jackson.databind.annotation.JsonDeserialize;

import java.util.List;
import java.util.Objects;

import static org.apache.paimon.utils.Preconditions.checkArgument;
import static org.apache.paimon.utils.Preconditions.checkNotNull;

/** TRIM/LTRIM/RTRIM {@link Transform}. */
public class TrimTransform extends StringTransform {
Expand All @@ -34,24 +41,43 @@ public class TrimTransform extends StringTransform {

private final Flag trimFlag;

public TrimTransform(List<Object> inputs, Flag trimFlag) {
public static final String FIELD_TRIM_FLAG = "trimFlag";

@JsonCreator
public TrimTransform(
@JsonProperty(StringTransform.FIELD_INPUTS)
@JsonDeserialize(contentUsing = StringTransform.InputDeserializer.class)
List<Object> inputs,
@JsonProperty(FIELD_TRIM_FLAG) Flag trimFlag) {
super(inputs);
this.trimFlag = trimFlag;
checkArgument(inputs.size() == 1 || inputs.size() == 2);
this.trimFlag = checkNotNull(trimFlag, "trimFlag must not be null");
}

@Override
public String name() {
return NAME;
}

@JsonGetter(FIELD_TRIM_FLAG)
public Flag trimFlag() {
return trimFlag;
}

@Override
public BinaryString transform(List<BinaryString> inputs) {
if (inputs.get(0) == null) {
return null;
}
String sourceString = inputs.get(0).toString();
String charsToTrim = inputs.size() == 1 ? " " : inputs.get(1).toString();
String charsToTrim = " ";
if (inputs.size() == 2) {
if (inputs.get(1) == null) {
// StringUtils.ltrim/rtrim treat a null charsToTrim as a null result
return null;
}
charsToTrim = inputs.get(1).toString();
}
switch (trimFlag) {
case BOTH:
return BinaryString.fromString(StringUtils.trim(sourceString, charsToTrim));
Expand All @@ -69,6 +95,20 @@ public Transform copyWithNewInputs(List<Object> inputs) {
return new TrimTransform(inputs, this.trimFlag);
}

@Override
public boolean equals(Object o) {
if (!super.equals(o)) {
return false;
}
TrimTransform that = (TrimTransform) o;
return trimFlag == that.trimFlag;
}

@Override
public int hashCode() {
return Objects.hash(super.hashCode(), trimFlag);
}

/** Enum of trim functions. */
public enum Flag {
LEADING,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -72,4 +72,13 @@ public void testConcatHybridInputs() {
BinaryString.fromString("-he")));
assertThat(result).isEqualTo(BinaryString.fromString("ha-he"));
}

@Test
public void testToStringWithNullInput() {
List<Object> inputs = new ArrayList<>();
inputs.add(BinaryString.fromString("a"));
inputs.add(null);

assertThat(new ConcatTransform(inputs).toString()).isEqualTo("CONCAT(a, null)");
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,41 @@ public void testSubstringRefInputs() {
assertThat(result).isEqualTo(BinaryString.fromString("ell"));
}

@Test
public void testNullPositionFieldYieldsNull() {
List<Object> inputs = new ArrayList<>();
inputs.add(new FieldRef(0, "f0", DataTypes.STRING()));
inputs.add(new FieldRef(1, "f1", DataTypes.INT()));
assertThat(
new SubstringTransform(inputs)
.transform(
GenericRow.of(
BinaryString.fromString("123-45-6789"), null)))
.isNull();

inputs.add(new FieldRef(2, "f2", DataTypes.INT()));
assertThat(
new SubstringTransform(inputs)
.transform(
GenericRow.of(
BinaryString.fromString("123-45-6789"), 8, null)))
.isNull();
}

@Test
public void testPositionsCountUtf16CodeUnits() {
List<Object> inputs = new ArrayList<>();
inputs.add(new FieldRef(0, "f0", DataTypes.STRING()));
inputs.add(2);
inputs.add(2);

Object result =
new SubstringTransform(inputs)
.transform(GenericRow.of(BinaryString.fromString("😀abc")));

assertThat(result).isEqualTo(BinaryString.fromString("?a"));
}

@Test
public void testSubstringRefInputUsesSourceFieldNullability() {
List<Object> inputs = new ArrayList<>();
Expand Down
Loading
Loading