nsivabalan commented on code in PR #18984:
URL: https://github.com/apache/hudi/pull/18984#discussion_r3616420755
##########
hudi-sync/hudi-hive-sync/src/main/java/org/apache/hudi/hive/HiveSyncConfigHolder.java:
##########
@@ -122,6 +122,23 @@ public class HiveSyncConfigHolder {
.defaultValue(1000)
.markAdvanced()
.withDocumentation("The number of partitions one batch when synchronous
partitions to hive.");
+ public static final ConfigProperty<Boolean> HIVE_SYNC_BATCHING_ENABLED =
ConfigProperty
+ .key("hoodie.datasource.hive_sync.batching.enabled")
+ .defaultValue(false)
+ .markAdvanced()
+ .sinceVersion("1.1.0")
Review Comment:
Fixed in the next push — both configs now use `sinceVersion("1.3.0")`.
##########
hudi-sync/hudi-hive-sync/src/main/java/org/apache/hudi/hive/HiveSyncConfigHolder.java:
##########
@@ -122,6 +122,23 @@ public class HiveSyncConfigHolder {
.defaultValue(1000)
.markAdvanced()
.withDocumentation("The number of partitions one batch when synchronous
partitions to hive.");
+ public static final ConfigProperty<Boolean> HIVE_SYNC_BATCHING_ENABLED =
ConfigProperty
+ .key("hoodie.datasource.hive_sync.batching.enabled")
+ .defaultValue(false)
+ .markAdvanced()
+ .sinceVersion("1.1.0")
+ .withDocumentation("When true, partition operations
(add/update/touch/drop) in the active sync mode "
Review Comment:
Fixed in the next push. Reworked the `HIVE_SYNC_BATCHING_ENABLED` doc to
scope precisely: HiveQL only, TOUCH is batched into `batch_num` chunks and
dispatched in parallel, SET_LOCATION dispatched in parallel (one statement per
partition), ADD-batching unchanged, DROP still serial, HMS/JDBC modes not
affected. Also tightened the `HIVE_SYNC_BATCHING_THREADS` doc.
##########
hudi-sync/hudi-hive-sync/src/main/java/org/apache/hudi/hive/util/HiveDriverPool.java:
##########
@@ -0,0 +1,340 @@
+/*
+ * 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.hudi.hive.util;
+
+import org.apache.hudi.exception.HoodieException;
+import org.apache.hudi.hive.HiveSyncConfig;
+import org.apache.hudi.hive.HoodieHiveSyncException;
+
+import org.apache.hadoop.hive.conf.HiveConf;
+import org.apache.hadoop.hive.ql.Driver;
+import org.apache.hadoop.hive.ql.session.SessionState;
+import org.apache.hadoop.security.UserGroupInformation;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.CancellationException;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.ThreadFactory;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import static
org.apache.hudi.sync.common.HoodieSyncConfig.META_SYNC_DATABASE_NAME;
+
+/**
+ * Pool of Hive {@link Driver} + {@link SessionState} pairs for parallel
HiveQL DDL.
+ *
+ * <p>Hive's {@code SessionState.start(state)} binds state to the calling
thread's
+ * thread-local, and {@code Driver} reads from that thread-local during {@code
run()}.
+ * A Driver constructed on one thread cannot be safely used from another. This
pool
+ * solves that by giving each slot its own dedicated worker thread (a
single-thread
+ * executor) — the Driver and SessionState are built on that thread by a
bootstrap
+ * task, and all subsequent SQL for that slot runs on the same thread.
+ *
+ * <p><b>Usage contract:</b> use this pool only for partition-row DDL
statements that
+ * are independent of each other and freely shuffleable across workers.
Table-level
+ * statements (createTable, schema evolution, USE database) must continue to
run on
+ * the session {@code Driver} held by {@code HiveQueryDDLExecutor} on the sync
driver
+ * thread. The pool is gated behind {@code
hoodie.datasource.hive_sync.batching.enabled}
+ * and is constructed only for HiveQL sync mode.
+ */
+public class HiveDriverPool implements AutoCloseable {
+
+ private static final Logger LOG =
LoggerFactory.getLogger(HiveDriverPool.class);
+
+ // Per-worker Driver construction has to be fast in practice (a few hundred
ms
+ // for the SessionState + Driver init). A 60s ceiling per worker leaves
plenty of
+ // headroom for a slow JVM warm-up but bounds the failure mode if the
metastore
+ // is unreachable or Hive hangs during init.
+ private static final long BOOTSTRAP_TIMEOUT_SECONDS = 60;
+
+ private final List<Worker> workers;
+ private final int size;
+ private volatile boolean closed;
+
+ public HiveDriverPool(HiveSyncConfig config, int size) {
+ this(config, size, new DefaultDriverFactory(config));
+ }
+
+ // Package-private for tests: accepts a DriverFactory so unit tests can
inject
+ // mock Driver instances without standing up a real Hive instance.
+ HiveDriverPool(HiveSyncConfig config, int size, DriverFactory factory) {
+ if (size < 1) {
+ throw new IllegalArgumentException("Pool size must be >= 1, got " +
size);
+ }
+ this.size = size;
+ this.workers = new ArrayList<>(size);
+ String databaseName = config.getStringOrDefault(META_SYNC_DATABASE_NAME);
+ PoolThreadFactory threadFactory = new PoolThreadFactory();
+ List<Future<Void>> bootstrapFutures = new ArrayList<>(size);
+ try {
+ for (int i = 0; i < size; i++) {
+ Worker worker = new Worker(threadFactory);
+ workers.add(worker);
+ bootstrapFutures.add(worker.executor.submit(() -> {
+ worker.driver = factory.newDriver(databaseName);
+ return null;
+ }));
+ }
+ // Block until all bootstraps complete so we surface construction errors
+ // before any caller hands us SQL. Bounded by BOOTSTRAP_TIMEOUT_SECONDS
so a
+ // hung Hive init doesn't deadlock the sync driver thread.
+ for (Future<Void> f : bootstrapFutures) {
+ f.get(BOOTSTRAP_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+ }
+ } catch (Exception e) {
+ tearDown();
+ throw new HoodieException("Failed to construct HiveDriverPool of size "
+ size, e);
+ }
+ LOG.info("Initialized HiveDriverPool with {} workers", size);
+ }
+
+ /**
+ * Runs each given SQL on <i>every</i> worker, in order. Used for setup
statements
+ * (e.g. {@code USE database}) that must establish per-thread session context
+ * before any partition statement runs. Blocks until all workers have
completed
+ * the setup. Throws on first error.
+ */
+ public void runOnEachWorker(List<String> setupSqls) {
+ if (closed) {
+ throw new IllegalStateException("Cannot dispatch to a closed
HiveDriverPool");
+ }
+ if (setupSqls.isEmpty()) {
+ return;
+ }
+ List<Future<?>> futures = new ArrayList<>(workers.size());
+ for (Worker worker : workers) {
+ futures.add(worker.executor.submit(() -> {
+ for (String sql : setupSqls) {
+ worker.driver.run(sql);
+ }
+ return null;
+ }));
+ }
+ awaitAll(futures);
+ }
+
+ /**
+ * Dispatches each SQL string to a worker (round-robin) and returns the list
of
+ * futures. The caller is responsible for awaiting and collecting errors.
SQL text
+ * is intentionally not logged per-statement here: batched TOUCH/ADD
statements can
+ * be many kilobytes, and N parallel workers would multiply the log volume.
See
+ * {@link #awaitAll(List)} for the per-call summary log.
+ */
+ public List<Future<?>> runAll(List<String> sqls) {
Review Comment:
Renamed to `dispatchAll` in the next push (and updated the Javadoc to
explicitly call out that it does not block — callers must `awaitAll`).
##########
hudi-sync/hudi-hive-sync/src/main/java/org/apache/hudi/hive/ddl/HiveQueryDDLExecutor.java:
##########
@@ -82,19 +95,61 @@ public void runSQL(String sql) {
updateHiveSQLs(Collections.singletonList(sql));
}
+ /**
+ * Partition-phase SQL fan-out. When the driver pool is present, any leading
+ * {@code USE database} statements are run on every worker (Hive 2.x's
+ * ALTER PARTITION SET LOCATION ignores db.table qualifiers and uses the
+ * connection's current database, so each worker needs to USE the right db
+ * before any partition ALTER). The remaining statements are then dispatched
+ * round-robin across the pool. Falls through to the sequential path on the
+ * session Driver when no pool is configured.
+ */
+ @Override
+ protected void runSQLs(List<String> sqls) {
+ if (sqls.isEmpty()) {
+ return;
+ }
+ if (!driverPool.isPresent()) {
+ updateHiveSQLs(sqls);
+ return;
+ }
+ HiveDriverPool pool = driverPool.get();
+ int firstNonUse = 0;
+ while (firstNonUse < sqls.size() && isUseStatement(sqls.get(firstNonUse)))
{
Review Comment:
Renamed `firstNonUse` → `useStatementCount` in the next push.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]