This is an automated email from the ASF dual-hosted git repository.
mchades pushed a commit to branch branch-1.3
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/branch-1.3 by this push:
new 115ad5a439 [Cherry-pick to branch-1.3] [#9374] fix(spark): Fix MySQL
datetime filter pushdown in Spark JDBC connector (#11451) (#11508)
115ad5a439 is described below
commit 115ad5a4396e06bef54db497576b014e281deb00
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Tue Jun 9 11:47:00 2026 +0800
[Cherry-pick to branch-1.3] [#9374] fix(spark): Fix MySQL datetime filter
pushdown in Spark JDBC connector (#11451) (#11508)
**Cherry-pick Information:**
- Original commit: eb2e2c9091701f6d6a2023c6cebd9385c5f01389
- Target branch: `branch-1.3`
- Status: ✅ Clean cherry-pick (no conflicts)
Co-authored-by: Yuhui <[email protected]>
---
.../test/jdbc/SparkJdbcMysqlCatalogIT.java | 55 ++++++++++++++++++++--
.../connector/jdbc/SparkJdbcTypeConverter34.java | 14 ++++--
.../jdbc/GravitinoJdbcCatalogSpark35.java | 3 +-
3 files changed, 62 insertions(+), 10 deletions(-)
diff --git
a/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/integration/test/jdbc/SparkJdbcMysqlCatalogIT.java
b/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/integration/test/jdbc/SparkJdbcMysqlCatalogIT.java
index f01e7807f5..fcb39806d9 100644
---
a/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/integration/test/jdbc/SparkJdbcMysqlCatalogIT.java
+++
b/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/integration/test/jdbc/SparkJdbcMysqlCatalogIT.java
@@ -9,10 +9,10 @@
*
* 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
+ * 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.
*/
@@ -21,6 +21,10 @@ package
org.apache.gravitino.spark.connector.integration.test.jdbc;
import static
org.apache.gravitino.integration.test.util.TestDatabaseName.MYSQL_CATALOG_MYSQL_IT;
import com.google.common.collect.Maps;
+import java.sql.Connection;
+import java.sql.DriverManager;
+import java.sql.Statement;
+import java.util.List;
import java.util.Map;
import org.apache.gravitino.credential.CredentialConstants;
import org.apache.gravitino.credential.JdbcCredential;
@@ -28,7 +32,9 @@ import
org.apache.gravitino.integration.test.container.ContainerSuite;
import org.apache.gravitino.spark.connector.integration.test.SparkCommonIT;
import
org.apache.gravitino.spark.connector.integration.test.util.SparkTableInfoChecker;
import org.apache.gravitino.spark.connector.jdbc.JdbcPropertiesConstants;
+import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
@Tag("gravitino-docker-test")
public abstract class SparkJdbcMysqlCatalogIT extends SparkCommonIT {
@@ -114,4 +120,45 @@ public abstract class SparkJdbcMysqlCatalogIT extends
SparkCommonIT {
CredentialConstants.CREDENTIAL_PROVIDERS,
JdbcCredential.JDBC_CREDENTIAL_TYPE);
return catalogProperties;
}
+
+ @Test
+ void testDatetimeFilterPushdown() throws Exception {
+ // Regression test for https://github.com/apache/gravitino/issues/9374:
+ // MySQL datetime columns must be usable as filter predicates via the
Gravitino Spark JDBC
+ // catalog. Previously, datetime was mapped to TimestampNTZType (Spark
3.4+), causing the
+ // MySQL JDBC dialect to generate invalid SQL literals during filter
pushdown.
+ //
+ // The table is created directly via JDBC to use MySQL-native datetime
type, exactly
+ // reproducing the original bug scenario where tables exist with datetime
columns in MySQL.
+ //
+ // Note: Uses subquery filter to avoid timezone conversion issues between
JVM (UTC+8) and
+ // MySQL container (UTC). The subquery approach validates filter pushdown
functionality
+ // while being resilient to timezone configuration differences in test
environments.
+ String db = getDefaultDatabase();
+ String tableName = "datetime_filter_test";
+ String jdbcUrl = mysqlUrl + "/" + db;
+ try (Connection conn = DriverManager.getConnection(jdbcUrl, mysqlUsername,
mysqlPassword);
+ Statement stmt = conn.createStatement()) {
+ stmt.execute("DROP TABLE IF EXISTS " + tableName);
+ stmt.execute("CREATE TABLE " + tableName + " (id INT, create_time
DATETIME)");
+ stmt.execute(
+ "INSERT INTO "
+ + tableName
+ + " VALUES"
+ + " (1, '2025-12-01 08:00:00'),"
+ + " (2, '2025-12-04 09:07:43'),"
+ + " (3, '2025-12-10 12:00:00')");
+ }
+
+ // Use subquery filter to avoid timezone-related assertion failures
+ // This validates filter pushdown works correctly after the TimestampType
fix
+ List<String> result =
+ getQueryData(
+ String.format(
+ "SELECT id FROM %s WHERE create_time >= (SELECT create_time
FROM %s WHERE id = 2) ORDER BY id",
+ tableName, tableName));
+ Assertions.assertEquals(2, result.size());
+ Assertions.assertTrue(result.contains("2"));
+ Assertions.assertTrue(result.contains("3"));
+ }
}
diff --git
a/spark-connector/v3.4/spark/src/main/java/org/apache/gravitino/spark/connector/jdbc/SparkJdbcTypeConverter34.java
b/spark-connector/v3.4/spark/src/main/java/org/apache/gravitino/spark/connector/jdbc/SparkJdbcTypeConverter34.java
index bbd32e0225..ae9659edcd 100644
---
a/spark-connector/v3.4/spark/src/main/java/org/apache/gravitino/spark/connector/jdbc/SparkJdbcTypeConverter34.java
+++
b/spark-connector/v3.4/spark/src/main/java/org/apache/gravitino/spark/connector/jdbc/SparkJdbcTypeConverter34.java
@@ -9,10 +9,10 @@
*
* 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
+ * 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.
*/
@@ -32,6 +32,12 @@ public class SparkJdbcTypeConverter34 extends
SparkTypeConverter34 {
// type varchar.
if (gravitinoType instanceof Types.VarCharType) {
return DataTypes.StringType;
+ } else if (gravitinoType instanceof Types.TimestampType
+ && !((Types.TimestampType) gravitinoType).hasTimeZone()) {
+ // MySQL datetime (no timezone) must map to TimestampType, not
TimestampNTZType, for JDBC
+ // filter pushdown compatibility. The MySQL JDBC dialect only handles
TimestampType literals
+ // correctly; TimestampNTZType produces invalid SQL syntax errors during
filter pushdown.
+ return DataTypes.TimestampType;
} else {
return super.toSparkType(gravitinoType);
}
diff --git
a/spark-connector/v3.5/spark/src/main/java/org/apache/gravitino/spark/connector/jdbc/GravitinoJdbcCatalogSpark35.java
b/spark-connector/v3.5/spark/src/main/java/org/apache/gravitino/spark/connector/jdbc/GravitinoJdbcCatalogSpark35.java
index 9ce36c2be1..9448664085 100644
---
a/spark-connector/v3.5/spark/src/main/java/org/apache/gravitino/spark/connector/jdbc/GravitinoJdbcCatalogSpark35.java
+++
b/spark-connector/v3.5/spark/src/main/java/org/apache/gravitino/spark/connector/jdbc/GravitinoJdbcCatalogSpark35.java
@@ -23,7 +23,6 @@ import java.util.Set;
import org.apache.gravitino.spark.connector.SparkTableChangeConverter;
import org.apache.gravitino.spark.connector.SparkTableChangeConverter34;
import org.apache.gravitino.spark.connector.SparkTypeConverter;
-import org.apache.gravitino.spark.connector.SparkTypeConverter34;
import org.apache.spark.sql.catalyst.analysis.NoSuchTableException;
import org.apache.spark.sql.connector.catalog.Identifier;
import org.apache.spark.sql.connector.catalog.Table;
@@ -33,7 +32,7 @@ public class GravitinoJdbcCatalogSpark35 extends
GravitinoJdbcCatalog {
@Override
protected SparkTypeConverter getSparkTypeConverter() {
- return new SparkTypeConverter34();
+ return new SparkJdbcTypeConverter34();
}
@Override