This is an automated email from the ASF dual-hosted git repository.
lgbo pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/incubator-gluten.git
The following commit(s) were added to refs/heads/main by this push:
new d164c8243 [GLUTEN-6935][CH]query fails when set session level
join_algorithm to grace_hash (#6944)
d164c8243 is described below
commit d164c8243a2027ee51df7a1661c92161520f234d
Author: loudongfeng <[email protected]>
AuthorDate: Wed Aug 21 10:42:20 2024 +0800
[GLUTEN-6935][CH]query fails when set session level join_algorithm to
grace_hash (#6944)
---
.../execution/GlutenClickHouseJoinSuite.scala | 106 +++++++++++++++++++++
cpp-ch/local-engine/Parser/JoinRelParser.cpp | 7 +-
2 files changed, 109 insertions(+), 4 deletions(-)
diff --git
a/backends-clickhouse/src/test/scala/org/apache/gluten/execution/GlutenClickHouseJoinSuite.scala
b/backends-clickhouse/src/test/scala/org/apache/gluten/execution/GlutenClickHouseJoinSuite.scala
new file mode 100644
index 000000000..75c4372a0
--- /dev/null
+++
b/backends-clickhouse/src/test/scala/org/apache/gluten/execution/GlutenClickHouseJoinSuite.scala
@@ -0,0 +1,106 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.gluten.execution
+
+import org.apache.gluten.GlutenConfig
+import org.apache.gluten.utils.UTSystemParameters
+
+import org.apache.spark.SparkConf
+
+class GlutenClickHouseJoinSuite extends
GlutenClickHouseWholeStageTransformerSuite {
+
+ protected val tablesPath: String = basePath + "/tpch-data"
+ protected val tpchQueries: String =
+ rootPath + "../../../../gluten-core/src/test/resources/tpch-queries"
+ protected val queriesResults: String = rootPath + "queries-output"
+
+ private val joinAlgorithm =
"spark.gluten.sql.columnar.backend.ch.runtime_settings.join_algorithm"
+
+ override protected def sparkConf: SparkConf = {
+ super.sparkConf
+ .set("spark.sql.files.maxPartitionBytes", "1g")
+ .set("spark.serializer", "org.apache.spark.serializer.JavaSerializer")
+ .set("spark.sql.shuffle.partitions", "5")
+ .set("spark.sql.adaptive.enabled", "false")
+ .set("spark.sql.files.minPartitionNum", "1")
+ .set("spark.gluten.sql.columnar.columnartorow", "true")
+ .set("spark.gluten.sql.columnar.backend.ch.worker.id", "1")
+ .set(GlutenConfig.GLUTEN_LIB_PATH, UTSystemParameters.clickHouseLibPath)
+ .set("spark.gluten.sql.columnar.iterator", "true")
+ .set("spark.gluten.sql.columnar.hashagg.enablefinal", "true")
+ .set("spark.gluten.sql.enable.native.validation", "false")
+ .set("spark.sql.warehouse.dir", warehouse)
+ .set(
+ "spark.sql.warehouse.dir",
+ getClass.getResource("/").getPath +
"tests-working-home/spark-warehouse")
+ .set("spark.hive.exec.dynamic.partition.mode", "nonstrict")
+ .set("spark.shuffle.manager", "sort")
+ .set("spark.io.compression.codec", "snappy")
+ .set("spark.sql.shuffle.partitions", "5")
+ .set("spark.sql.autoBroadcastJoinThreshold", "10MB")
+ .set(joinAlgorithm, "hash")
+ .set("spark.sql.autoBroadcastJoinThreshold", "-1")
+ .setMaster("local[*]")
+ }
+
+ test("int to long join key rewrite causes column miss match ") {
+ assert("hash".equalsIgnoreCase(sparkConf.get(joinAlgorithm, "hash")))
+ withSQLConf(joinAlgorithm -> "grace_hash") {
+ withTable("my_customer", "my_store_sales", "my_date_dim") {
+ sql("""
+ |CREATE TABLE my_customer (
+ | c_customer_sk INT)
+ |USING orc
+ |""".stripMargin)
+ sql("""
+ |CREATE TABLE my_store_sales (
+ | ss_sold_date_sk INT,
+ | ss_customer_sk INT)
+ | USING orc
+ |""".stripMargin)
+ sql("""
+ |CREATE TABLE my_date_dim (
+ | d_date_sk INT,
+ | d_year INT,
+ | d_qoy INT)
+ |USING orc
+ |""".stripMargin)
+
+ sql("insert into my_customer values (1), (2), (3), (4)")
+ sql("insert into my_store_sales values (1, 1), (2, 2), (3, 3), (4, 4)")
+ sql("insert into my_date_dim values (1, 2002, 1), (2, 2002, 2)")
+ val q =
+ """
+ |SELECT
+ | count(*) cnt1
+ |FROM
+ | my_customer c
+ |WHERE
+ | exists(SELECT *
+ | FROM my_store_sales, my_date_dim
+ | WHERE c.c_customer_sk = ss_customer_sk AND
+ | ss_sold_date_sk = d_date_sk AND
+ | d_year = 2002 AND
+ | d_qoy < 4)
+ |LIMIT 100
+ |""".stripMargin
+
runQueryAndCompare(q)(checkGlutenOperatorMatch[CHShuffledHashJoinExecTransformer])
+ }
+ }
+ }
+
+}
diff --git a/cpp-ch/local-engine/Parser/JoinRelParser.cpp
b/cpp-ch/local-engine/Parser/JoinRelParser.cpp
index 30651aff1..0446a397c 100644
--- a/cpp-ch/local-engine/Parser/JoinRelParser.cpp
+++ b/cpp-ch/local-engine/Parser/JoinRelParser.cpp
@@ -53,11 +53,10 @@ using namespace DB;
namespace local_engine
{
-std::shared_ptr<DB::TableJoin>
createDefaultTableJoin(substrait::JoinRel_JoinType join_type, bool
is_existence_join)
+std::shared_ptr<DB::TableJoin>
createDefaultTableJoin(substrait::JoinRel_JoinType join_type, bool
is_existence_join, ContextPtr & context)
{
- auto & global_context = SerializedPlanParser::global_context;
auto table_join = std::make_shared<TableJoin>(
- global_context->getSettingsRef(),
global_context->getGlobalTemporaryVolume(),
global_context->getTempDataOnDisk());
+ context->getSettingsRef(), context->getGlobalTemporaryVolume(),
context->getTempDataOnDisk());
std::pair<DB::JoinKind, DB::JoinStrictness> kind_and_strictness =
JoinUtil::getJoinKindAndStrictness(join_type, is_existence_join);
table_join->setKind(kind_and_strictness.first);
@@ -216,7 +215,7 @@ DB::QueryPlanPtr JoinRelParser::parseJoin(const
substrait::JoinRel & join, DB::Q
renamePlanColumns(*left, *right, *storage_join);
}
- auto table_join = createDefaultTableJoin(join.type(),
join_opt_info.is_existence_join);
+ auto table_join = createDefaultTableJoin(join.type(),
join_opt_info.is_existence_join, context);
DB::Block right_header_before_convert_step =
right->getCurrentDataStream().header;
addConvertStep(*table_join, *left, *right);
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]