yihua commented on code in PR #11790:
URL: https://github.com/apache/hudi/pull/11790#discussion_r1759787749


##########
hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/transaction/lock/ZookeeperBasedImplicitBasePathLockProvider.java:
##########
@@ -0,0 +1,70 @@
+/*
+ * 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.client.transaction.lock;
+
+import org.apache.hudi.common.config.HoodieCommonConfig;
+import org.apache.hudi.common.config.LockConfiguration;
+import org.apache.hudi.common.lock.LockProvider;
+import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.util.ValidationUtils;
+import org.apache.hudi.common.util.hash.HashID;
+import org.apache.hudi.storage.StorageConfiguration;
+
+import javax.annotation.concurrent.NotThreadSafe;
+
+import static org.apache.hudi.aws.utils.S3Utils.s3aToS3;
+import static org.apache.hudi.common.util.StringUtils.concatenateWithThreshold;
+
+/**
+ * A zookeeper based lock. This {@link LockProvider} implementation allows to 
lock table operations
+ * using zookeeper. Users need to have a Zookeeper cluster deployed to be able 
to use this lock.
+ *
+ * This class derives the zookeeper base path from the hudi table base path 
(hoodie.base.path) and
+ * table name (hoodie.table.name), with lock key set to a hard-coded value.
+ */
+@NotThreadSafe
+public class ZookeeperBasedImplicitBasePathLockProvider extends 
BaseZookeeperBasedLockProvider {
+
+  public static final String LOCK_KEY = "lock_key";
+
+  public static String getLockBasePath(String hudiTableBasePath, String 
hudiTableName) {
+    // Ensure consistent format for S3 URI.
+    String hashPart = '-' + 
HashID.generateXXHashAsString(s3aToS3(hudiTableBasePath), HashID.Size.BITS_64);
+    String folderName = concatenateWithThreshold(hudiTableName, hashPart, 
MAX_ZK_BASE_PATH_NUM_BYTES);
+    return "/tmp/" + folderName;
+  }
+
+  public ZookeeperBasedImplicitBasePathLockProvider(final LockConfiguration 
lockConfiguration, final StorageConfiguration<?> conf) {
+    super(lockConfiguration, conf);
+  }
+
+  @Override
+  protected String getZkBasePath(LockConfiguration lockConfiguration) {
+    String hudiTableBasePath = 
lockConfiguration.getConfig().getString(HoodieCommonConfig.BASE_PATH.key());

Review Comment:
   Use `ConfigUtils#getStringWithAltKeys`



##########
hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/transaction/lock/BaseZookeeperBasedLockProvider.java:
##########
@@ -0,0 +1,196 @@
+/*
+ * 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.client.transaction.lock;
+
+import org.apache.hudi.common.config.LockConfiguration;
+import org.apache.hudi.common.lock.LockProvider;
+import org.apache.hudi.common.lock.LockState;
+import org.apache.hudi.common.util.StringUtils;
+import org.apache.hudi.common.util.ValidationUtils;
+import org.apache.hudi.exception.HoodieLockException;
+import org.apache.hudi.storage.StorageConfiguration;
+
+import org.apache.curator.framework.CuratorFramework;
+import org.apache.curator.framework.CuratorFrameworkFactory;
+import org.apache.curator.framework.recipes.locks.InterProcessMutex;
+import org.apache.curator.retry.BoundedExponentialBackoffRetry;
+import org.apache.zookeeper.KeeperException;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import javax.annotation.concurrent.NotThreadSafe;
+
+import java.io.Serializable;
+import java.util.concurrent.TimeUnit;
+
+import static 
org.apache.hudi.common.config.LockConfiguration.DEFAULT_ZK_CONNECTION_TIMEOUT_MS;
+import static 
org.apache.hudi.common.config.LockConfiguration.DEFAULT_ZK_SESSION_TIMEOUT_MS;
+import static 
org.apache.hudi.common.config.LockConfiguration.LOCK_ACQUIRE_NUM_RETRIES_PROP_KEY;
+import static 
org.apache.hudi.common.config.LockConfiguration.LOCK_ACQUIRE_RETRY_MAX_WAIT_TIME_IN_MILLIS_PROP_KEY;
+import static 
org.apache.hudi.common.config.LockConfiguration.LOCK_ACQUIRE_RETRY_WAIT_TIME_IN_MILLIS_PROP_KEY;
+import static 
org.apache.hudi.common.config.LockConfiguration.ZK_CONNECTION_TIMEOUT_MS_PROP_KEY;
+import static 
org.apache.hudi.common.config.LockConfiguration.ZK_CONNECT_URL_PROP_KEY;
+import static 
org.apache.hudi.common.config.LockConfiguration.ZK_SESSION_TIMEOUT_MS_PROP_KEY;
+
+/**
+ * A zookeeper based lock. This {@link LockProvider} implementation allows to 
lock table operations
+ * using zookeeper. Users need to have a Zookeeper cluster deployed to be able 
to use this lock.
+ */
+@NotThreadSafe
+public abstract class BaseZookeeperBasedLockProvider implements 
LockProvider<InterProcessMutex>, Serializable {
+
+  private static final Logger LOG = 
LoggerFactory.getLogger(BaseZookeeperBasedLockProvider.class);
+
+  private final transient CuratorFramework curatorFrameworkClient;
+  private volatile InterProcessMutex lock = null;
+  protected final LockConfiguration lockConfiguration;
+  protected final String zkBasePath;
+  protected final String lockKey;
+
+  public static final int MAX_ZK_BASE_PATH_NUM_BYTES = 4096;
+
+  public BaseZookeeperBasedLockProvider(final LockConfiguration 
lockConfiguration, final StorageConfiguration<?> conf) {
+    checkRequiredProps(lockConfiguration);
+    this.lockConfiguration = lockConfiguration;
+    zkBasePath = getZkBasePath(lockConfiguration);
+    lockKey = getLockKey(lockConfiguration);
+    this.curatorFrameworkClient = CuratorFrameworkFactory.builder()
+        
.connectString(lockConfiguration.getConfig().getString(ZK_CONNECT_URL_PROP_KEY))
+        .retryPolicy(new 
BoundedExponentialBackoffRetry(lockConfiguration.getConfig().getInteger(LOCK_ACQUIRE_RETRY_WAIT_TIME_IN_MILLIS_PROP_KEY),
+            
lockConfiguration.getConfig().getInteger(LOCK_ACQUIRE_RETRY_MAX_WAIT_TIME_IN_MILLIS_PROP_KEY),
+            
lockConfiguration.getConfig().getInteger(LOCK_ACQUIRE_NUM_RETRIES_PROP_KEY)))
+        
.sessionTimeoutMs(lockConfiguration.getConfig().getInteger(ZK_SESSION_TIMEOUT_MS_PROP_KEY,
 DEFAULT_ZK_SESSION_TIMEOUT_MS))
+        
.connectionTimeoutMs(lockConfiguration.getConfig().getInteger(ZK_CONNECTION_TIMEOUT_MS_PROP_KEY,
 DEFAULT_ZK_CONNECTION_TIMEOUT_MS))
+        .build();
+    this.curatorFrameworkClient.start();
+    createPathIfNotExists();
+  }
+
+  protected abstract String getZkBasePath(LockConfiguration lockConfiguration);
+
+  protected abstract String getLockKey(LockConfiguration lockConfiguration);
+
+  protected String generateLogSuffixString() {
+    return StringUtils.join("ZkBasePath = ", zkBasePath, ", lock key = ", 
lockKey);
+  }
+
+  protected String getLockPath() {
+    return zkBasePath + '/' + lockKey;
+  }
+
+  private void createPathIfNotExists() {
+    try {
+      String lockPath = getLockPath();
+      LOG.info(String.format("Creating zookeeper path %s if not exists", 
lockPath));
+      String[] parts = lockPath.split("/");
+      StringBuilder currentPath = new StringBuilder();
+      for (String part : parts) {
+        if (!part.isEmpty()) {
+          currentPath.append("/").append(part);
+          createNodeIfNotExists(currentPath.toString());
+        }
+      }
+    } catch (Exception e) {
+      LOG.error("Failed to create ZooKeeper path: " + e.getMessage());
+      throw new HoodieLockException("Failed to initialize ZooKeeper path", e);
+    }
+  }
+
+  private void createNodeIfNotExists(String path) throws Exception {

Review Comment:
   Do we want to keep the same logic as before:
   ```
   private void createNodeIfNotExists(String path) throws Exception {
       if (this.curatorFrameworkClient.checkExists().forPath(path) == null) {
         try {
           this.curatorFrameworkClient.create().forPath(path);
           // to avoid failure due to synchronous calls.
         } catch (KeeperException e) {
           if (e.code() == KeeperException.Code.NODEEXISTS) {
             LOG.debug(String.format("Node already exist for path = %s", path));
           } else {
             throw new HoodieLockException("Failed to create zookeeper node", 
e);
           }
         }
       }
     }
   ```



-- 
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]

Reply via email to