This is an automated email from the ASF dual-hosted git repository.
hyuan pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/calcite.git
The following commit(s) were added to refs/heads/master by this push:
new 6218661 [CALCITE-3867] Support RelDistribution json serialization
(Krisztian Kasa)
6218661 is described below
commit 62186610f430996b28295d930a4ac505eed17b81
Author: kkasa <[email protected]>
AuthorDate: Mon Mar 9 16:20:31 2020 +0100
[CALCITE-3867] Support RelDistribution json serialization (Krisztian Kasa)
Close #1868
---
.../org/apache/calcite/rel/RelDistributions.java | 6 +-
.../apache/calcite/rel/externalize/RelJson.java | 36 ++++++++++-
.../calcite/rel/externalize/RelJsonReader.java | 2 +-
.../calcite/rel/logical/LogicalSortExchange.java | 8 +++
.../org/apache/calcite/plan/RelWriterTest.java | 69 ++++++++++++++++++++++
5 files changed, 117 insertions(+), 4 deletions(-)
diff --git a/core/src/main/java/org/apache/calcite/rel/RelDistributions.java
b/core/src/main/java/org/apache/calcite/rel/RelDistributions.java
index b5444d1..1de39fd 100644
--- a/core/src/main/java/org/apache/calcite/rel/RelDistributions.java
+++ b/core/src/main/java/org/apache/calcite/rel/RelDistributions.java
@@ -35,7 +35,7 @@ import javax.annotation.Nonnull;
* Utilities concerning {@link org.apache.calcite.rel.RelDistribution}.
*/
public class RelDistributions {
- private static final ImmutableIntList EMPTY = ImmutableIntList.of();
+ public static final ImmutableIntList EMPTY = ImmutableIntList.of();
/** The singleton singleton distribution. */
public static final RelDistribution SINGLETON =
@@ -80,6 +80,10 @@ public class RelDistributions {
return RelDistributionTraitDef.INSTANCE.canonize(trait);
}
+ public static RelDistribution of(RelDistribution.Type type, ImmutableIntList
keys) {
+ return new RelDistributionImpl(type, keys);
+ }
+
/** Implementation of {@link org.apache.calcite.rel.RelDistribution}. */
private static class RelDistributionImpl implements RelDistribution {
private static final Ordering<Iterable<Integer>> ORDERING =
diff --git a/core/src/main/java/org/apache/calcite/rel/externalize/RelJson.java
b/core/src/main/java/org/apache/calcite/rel/externalize/RelJson.java
index f6715bf..ea45bc1 100644
--- a/core/src/main/java/org/apache/calcite/rel/externalize/RelJson.java
+++ b/core/src/main/java/org/apache/calcite/rel/externalize/RelJson.java
@@ -58,6 +58,7 @@ import org.apache.calcite.sql.parser.SqlParserPos;
import org.apache.calcite.sql.type.SqlTypeName;
import org.apache.calcite.sql.validate.SqlNameMatchers;
import org.apache.calcite.util.ImmutableBitSet;
+import org.apache.calcite.util.ImmutableIntList;
import org.apache.calcite.util.JsonBuilder;
import org.apache.calcite.util.Util;
@@ -73,6 +74,8 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
+import static org.apache.calcite.rel.RelDistributions.EMPTY;
+
/**
* Utilities for converting {@link org.apache.calcite.rel.RelNode}
* into JSON format.
@@ -190,8 +193,35 @@ public class RelJson {
return new RelFieldCollation(field, direction, nullDirection);
}
- public RelDistribution toDistribution(Object o) {
- return RelDistributions.ANY; // TODO:
+ public RelDistribution toDistribution(Map<String, Object> map) {
+ final RelDistribution.Type type =
+ Util.enumVal(RelDistribution.Type.class,
+ (String) map.get("type"));
+
+ ImmutableIntList list = EMPTY;
+ if (map.containsKey("keys")) {
+ List<Object> keysJson = (List<Object>) map.get("keys");
+ ArrayList<Integer> keys = new ArrayList<>(keysJson.size());
+ for (Object o : keysJson) {
+ keys.add((Integer) o);
+ }
+ list = ImmutableIntList.copyOf(keys);
+ }
+ return RelDistributions.of(type, list);
+ }
+
+ private Object toJson(RelDistribution relDistribution) {
+ final Map<String, Object> map = jsonBuilder.map();
+ map.put("type", relDistribution.getType().name());
+
+ if (!relDistribution.getKeys().isEmpty()) {
+ List<Object> keys = new ArrayList<>(relDistribution.getKeys().size());
+ for (Integer key : relDistribution.getKeys()) {
+ keys.add(toJson(key));
+ }
+ map.put("keys", keys);
+ }
+ return map;
}
public RelDataType toType(RelDataTypeFactory typeFactory, Object o) {
@@ -285,6 +315,8 @@ public class RelJson {
return toJson((RelDataType) value);
} else if (value instanceof RelDataTypeField) {
return toJson((RelDataTypeField) value);
+ } else if (value instanceof RelDistribution) {
+ return toJson((RelDistribution) value);
} else {
throw new UnsupportedOperationException("type not serializable: "
+ value + " (type " + value.getClass().getCanonicalName() + ")");
diff --git
a/core/src/main/java/org/apache/calcite/rel/externalize/RelJsonReader.java
b/core/src/main/java/org/apache/calcite/rel/externalize/RelJsonReader.java
index 5f056a0..ce7c37a 100644
--- a/core/src/main/java/org/apache/calcite/rel/externalize/RelJsonReader.java
+++ b/core/src/main/java/org/apache/calcite/rel/externalize/RelJsonReader.java
@@ -234,7 +234,7 @@ public class RelJsonReader {
}
public RelDistribution getDistribution() {
- return relJson.toDistribution(get("distribution"));
+ return relJson.toDistribution((Map<String, Object>)
get("distribution"));
}
public ImmutableList<ImmutableList<RexLiteral>> getTuples(String tag) {
diff --git
a/core/src/main/java/org/apache/calcite/rel/logical/LogicalSortExchange.java
b/core/src/main/java/org/apache/calcite/rel/logical/LogicalSortExchange.java
index d1e6b5b..3870d5b 100644
--- a/core/src/main/java/org/apache/calcite/rel/logical/LogicalSortExchange.java
+++ b/core/src/main/java/org/apache/calcite/rel/logical/LogicalSortExchange.java
@@ -23,6 +23,7 @@ import org.apache.calcite.rel.RelCollation;
import org.apache.calcite.rel.RelCollationTraitDef;
import org.apache.calcite.rel.RelDistribution;
import org.apache.calcite.rel.RelDistributionTraitDef;
+import org.apache.calcite.rel.RelInput;
import org.apache.calcite.rel.RelNode;
import org.apache.calcite.rel.core.SortExchange;
@@ -37,6 +38,13 @@ public class LogicalSortExchange extends SortExchange {
}
/**
+ * Creates a LogicalSortExchange by parsing serialized output.
+ */
+ public LogicalSortExchange(RelInput input) {
+ super(input);
+ }
+
+ /**
* Creates a LogicalSortExchange.
*
* @param input Input relational expression
diff --git a/core/src/test/java/org/apache/calcite/plan/RelWriterTest.java
b/core/src/test/java/org/apache/calcite/plan/RelWriterTest.java
index 1714f2e..aab731a 100644
--- a/core/src/test/java/org/apache/calcite/plan/RelWriterTest.java
+++ b/core/src/test/java/org/apache/calcite/plan/RelWriterTest.java
@@ -19,6 +19,8 @@ package org.apache.calcite.plan;
import org.apache.calcite.adapter.java.ReflectiveSchema;
import org.apache.calcite.avatica.util.TimeUnit;
import org.apache.calcite.rel.RelCollations;
+import org.apache.calcite.rel.RelDistribution;
+import org.apache.calcite.rel.RelDistributions;
import org.apache.calcite.rel.RelNode;
import org.apache.calcite.rel.RelShuttleImpl;
import org.apache.calcite.rel.core.AggregateCall;
@@ -62,6 +64,7 @@ import org.apache.calcite.util.TestUtil;
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableSet;
+import com.google.common.collect.Lists;
import org.junit.jupiter.api.Test;
@@ -349,6 +352,37 @@ public class RelWriterTest {
+ " ]\n"
+ "}";
+ public static final String XX3 = "{\n"
+ + " \"rels\": [\n"
+ + " {\n"
+ + " \"id\": \"0\",\n"
+ + " \"relOp\": \"LogicalTableScan\",\n"
+ + " \"table\": [\n"
+ + " \"scott\",\n"
+ + " \"EMP\"\n"
+ + " ],\n"
+ + " \"inputs\": []\n"
+ + " },\n"
+ + " {\n"
+ + " \"id\": \"1\",\n"
+ + " \"relOp\": \"LogicalSortExchange\",\n"
+ + " \"distribution\": {\n"
+ + " \"type\": \"HASH_DISTRIBUTED\",\n"
+ + " \"keys\": [\n"
+ + " 0\n"
+ + " ]\n"
+ + " },\n"
+ + " \"collation\": [\n"
+ + " {\n"
+ + " \"field\": 0,\n"
+ + " \"direction\": \"ASCENDING\",\n"
+ + " \"nulls\": \"LAST\"\n"
+ + " }\n"
+ + " ]\n"
+ + " }\n"
+ + " ]\n"
+ + "}";
+
/**
* Unit test for {@link org.apache.calcite.rel.externalize.RelJsonWriter} on
* a simple tree of relational expressions, consisting of a table and a
@@ -843,4 +877,39 @@ public class RelWriterTest {
.build();
return rel;
}
+
+ @Test public void testWriteSortExchangeWithHashDistribution() {
+ final RelNode root =
createSortPlan(RelDistributions.hash(Lists.newArrayList(0)));
+ final RelJsonWriter writer = new RelJsonWriter();
+ root.explain(writer);
+ final String json = writer.asString();
+ assertThat(json, is(XX3));
+
+ final String s = deserializeAndDumpToTextFormat(getSchema(root), json);
+ final String expected =
+ "LogicalSortExchange(distribution=[hash[0]], collation=[[0]])\n"
+ + " LogicalTableScan(table=[[scott, EMP]])\n";
+ assertThat(s, isLinux(expected));
+ }
+
+ @Test public void testWriteSortExchangeWithRandomDistribution() {
+ final RelNode root = createSortPlan(RelDistributions.RANDOM_DISTRIBUTED);
+ final RelJsonWriter writer = new RelJsonWriter();
+ root.explain(writer);
+ final String json = writer.asString();
+ final String s = deserializeAndDumpToTextFormat(getSchema(root), json);
+ final String expected =
+ "LogicalSortExchange(distribution=[random], collation=[[0]])\n"
+ + " LogicalTableScan(table=[[scott, EMP]])\n";
+ assertThat(s, isLinux(expected));
+ }
+
+ private RelNode createSortPlan(RelDistribution distribution) {
+ final FrameworkConfig config = RelBuilderTest.config().build();
+ final RelBuilder builder = RelBuilder.create(config);
+ return builder.scan("EMP")
+ .sortExchange(distribution,
+ RelCollations.of(0))
+ .build();
+ }
}