•34 min read

分散システムにおける分散ロック: Redlock、PostgreSQLアドバイザリロック、etcdリース

分散システムにおける分散ロック: Redlock、PostgreSQLアドバイザリロック、etcdリース

分散ロックは、分散システムにおける共有リソースへの同時アクセスを調整するための基本的なプリミティブです。その正しい実装は、データの一貫性を保ち、競合状態(race condition)を防ぐ上で極めて重要です。このガイドでは、一般的な分散ロックのパターン、その障害モード、そしてRedlock、PostgreSQLアドバイザリロック、etcdリースを使った実践的な実装について詳しく解説します。

Audio Briefing
0:00 / 0:00

分散ロックの課題

単一プロセス内のミューテックスやセマフォとは異なり、分散ロックはネットワーク境界を越えて動作するため、ネットワークパーティション、ノード障害、予測不能なメッセージ遅延といった複雑な問題を引き起こします。堅牢な分散ロックは、以下の3つの主要な特性を満たす必要があります。

  1. 相互排他(Mutual Exclusion): ある時点において、ロックを保持できるクライアントは最大1つであること。
  2. 活性(Liveness、デッドロックフリー): クライアントがロックを取得した後にクラッシュした場合でも、最終的には他のクライアントがロックを取得できること。
  3. 耐障害性(Fault Tolerance): ロックシステム自体が個々のノードの障害に対して回復力があること。

4番目の、しばしば見落とされがちな特性として**フェンシング(Fencing)**があります。フェンシングは、ロックを取得したと「思っていた」がその後取り消された古い、遅い、またはクラッシュしたクライアントが、新しい正当なロック保持者の操作を妨害しないように保証するものです。

Advertisement

Redlock: 複数インスタンスRedisアプローチ

Redisの作者であるSalvatore Sanfilippoによって提案されたRedlockは、独立したRedisインスタンスの過半数がロックを許可することを要求することで、耐障害性を実現することを目指しています。

Redlockアルゴリズムの概要

  1. クライアントは現在のタイムスタンプをミリ秒単位で取得します。
  2. クライアントはN/2 + 1個のRedisインスタンスで、SET key random_value NX PX expiry_timeを使用してロックを取得しようとします。NXはキーが存在しない場合にのみ設定されることを保証し、PX expiry_timeは有効期限を設定します。random_valueはクライアントの一意のトークンです。
  3. クライアントはステップ1からの経過時間を計算します。
  4. ロックが過半数のインスタンスで取得され、かつ経過時間がロックの有効期間(expiry_time - time_elapsed)よりも短い場合、ロックは取得されたと見なされます。
  5. ロックが取得されなかった場合(インスタンスが十分でないか、有効期間が過ぎた場合)、クライアントは触れたすべてのインスタンスでロックを解放しようとします。

批判とフェンシングトークン

Martin KleppmannによるRedlockの批判は、重大な障害モードを浮き彫りにしています。

  • クロックドリフト(Clock Drift): Redisインスタンスまたはクライアントのクロックが大幅にドリフトすると、expiry_timeが誤って解釈され、複数のクライアントがロックを保持していると信じてしまう可能性があります。
  • GCポーズ/プロセス停止(GC Pauses/Process Stalls): ロックを保持しているクライアントが、長いガベージコレクションの一時停止やプロセス停止を経験する可能性があります。この停止中に、そのロックが期限切れになることがあります。クライアントが再開したとき、別のクライアントがすでにロックを取得して操作を実行していることに気づかずに、共有リソースに対して操作を行う可能性があります。これは相互排他を侵害します。
  • ネットワークパーティション(Network Partitions): クライアントがロックを取得した後、Redisインスタンスの過半数からパーティション化される可能性があります。そのクライアントは動作を続けるかもしれませんが、その間に別のクライアントが新しい過半数からロックを取得する可能性があります。

核心的な問題は、Redlock自体がフェンシングを提供しないことです。フェンシングトークンは、ロックサービスがロックを付与するたびに発行される単調増加する数値です。クライアントがロックを取得すると、このトークンを受け取ります。共有リソースに対するすべての操作にはこのトークンを含める必要があり、リソース自体はトークンが最新のものであることを検証する必要があります。

Redlockでのフェンシングの実装(概念)

Redlockはネイティブでフェンシングトークンをサポートしていませんが、堅牢な実装には、外部の強力に一貫性のあるカウンター、またはそのようなトークンを提供するためのRedisインスタンスへの変更が必要になります。これはしばしば、より複雑な合意形成ベースのシステムへと向かわせます。

// Simplified Redlock client (conceptual, not production-ready without fencing)
import { Redis } from 'ioredis';
import { v4 as uuidv4 } from 'uuid';

interface RedlockOptions {
  lockKey: string;
  resourceId: string;
  ttlMs: number; // Time-to-live for the lock in milliseconds
  retryDelayMs?: number;
  maxRetries?: number;
}

class RedlockClient {
  private redisClients: Redis[];
  private quorum: number;

  constructor(redisUrls: string[]) {
    this.redisClients = redisUrls.map(url => new Redis(url));
    this.quorum = Math.floor(redisUrls.length / 2) + 1;
    if (this.quorum === 0) {
      throw new Error("Redlock requires at least one Redis instance.");
    }
  }

  /**
   * Attempts to acquire a distributed lock using the Redlock algorithm.
   * @returns A unique token if the lock is acquired, otherwise null.
   */
  public async acquireLock(options: RedlockOptions): Promise<string | null> {
    const { lockKey, resourceId, ttlMs, retryDelayMs = 100, maxRetries = 5 } = options;
    const value = uuidv4(); // Unique token for this lock attempt
    let retries = 0;

    while (retries < maxRetries) {
      const startTime = Date.now();
      let acquiredCount = 0;

      const acquirePromises = this.redisClients.map(async (client) => {
        try {
          // SET key value NX PX ttlMs
          // NX: Only set if the key does not already exist.
          // PX: Set the specified expire time, in milliseconds.
          const result = await client.set(lockKey, value, 'NX', 'PX', ttlMs);
          return result === 'OK';
        } catch (error) {
          console.error(`Redis client error during acquire: ${error}`);
          return false;
        }
      });

      const results = await Promise.all(acquirePromises);
      acquiredCount = results.filter(Boolean).length;

      const elapsedTime = Date.now() - startTime;
      const isValid = elapsedTime < ttlMs;

      if (acquiredCount >= this.quorum && isValid) {
        console.log(`Lock '${lockKey}' acquired by '${value}' on ${acquiredCount} instances.`);
        return value; // Lock acquired successfully
      } else {
        // If not acquired, or validity time expired, release any locks we did get
        await this.releaseLock(lockKey, value);
        console.warn(`Failed to acquire lock '${lockKey}'. Acquired ${acquiredCount}/${this.quorum} instances. Retrying...`);
        retries++;
        await new Promise(resolve => setTimeout(resolve, retryDelayMs));
      }
    }

    console.error(`Failed to acquire lock '${lockKey}' after ${maxRetries} retries.`);
    return null;
  }

  /**
   * Releases a distributed lock.
   * It's crucial to only release locks that match the client's unique value.
   */
  public async releaseLock(lockKey: string, value: string): Promise<void> {
    // Lua script to ensure atomic check-and-delete
    const luaScript = `
      if redis.call("get", KEYS[1]) == ARGV[1] then
        return redis.call("del", KEYS[1])
      else
        return 0
      end
    `;

    const releasePromises = this.redisClients.map(async (client) => {
      try {
        await client.eval(luaScript, 1, lockKey, value);
      } catch (error) {
        console.error(`Redis client error during release: ${error}`);
      }
    });

    await Promise.all(releasePromises);
    console.log(`Lock '${lockKey}' released by '${value}'.`);
  }

  public async disconnect(): Promise<void> {
    await Promise.all(this.redisClients.map(client => client.quit()));
  }
}

// Example Usage (requires running Redis instances)
async function runRedlockExample() {
  const redisUrls = ['redis://localhost:6379', 'redis://localhost:6380', 'redis://localhost:6381'];
  const redlock = new RedlockClient(redisUrls);

  const lockKey = 'my_critical_resource';
  const resourceId = 'unique_resource_id_123';
  const ttlMs = 5000; // 5 seconds

  console.log('Attempting to acquire lock 1...');
  const token1 = await redlock.acquireLock({ lockKey, resourceId, ttlMs });

  if (token1) {
    console.log(`Client 1 acquired lock with token: ${token1}. Performing critical operation...`);
    // Simulate work
    await new Promise(resolve => setTimeout(resolve, 2000));
    console.log('Client 1 finished critical operation. Releasing lock...');
    await redlock.releaseLock(lockKey, token1);
  } else {
    console.log('Client 1 failed to acquire lock.');
  }

  console.log('\nAttempting to acquire lock 2 (should succeed after lock 1 is released)...');
  const token2 = await redlock.acquireLock({ lockKey, resourceId, ttlMs });

  if (token2) {
    console.log(`Client 2 acquired lock with token: ${token2}. Performing critical operation...`);
    await new Promise(resolve => setTimeout(resolve, 1000));
    await redlock.releaseLock(lockKey, token2);
  } else {
    console.log('Client 2 failed to acquire lock.');
  }

  await redlock.disconnect();
}

// To run this, you'd need 3 Redis instances running, e.g.:
// docker run -p 6379:6379 --name redis1 -d redis
// docker run -p 6380:6379 --name redis2 -d redis
// docker run -p 6381:6379 --name redis3 -d redis
// Then: node your_script.js
// runRedlockExample();

PostgreSQLアドバイザリロック

すでにPostgreSQLに大きく依存しているシステムにとって、アドバイザリロックは、よりシンプルでトランザクション対応の、非常に効率的な分散ロックメカニズムを提供します。これらは「アドバイザリ」と呼ばれます。なぜなら、PostgreSQLはその使用を強制しないため、一貫してロックを取得および解放するのはアプリケーションの責任だからです。

主な特徴

  • トランザクションスコープ(Transaction-Scoped): ロックはトランザクションの終了時(コミットまたはロールバック)に自動的に解放されます。これは活性を保証するための強力な機能です。
  • セッションスコープ(Session-Scoped): ロックはトランザクションとは独立して、セッションの期間中取得することもできます。
  • 軽量(Lightweight): テーブルアクセスをブロックしたり、行レベルロックを取得したりしないため、非常に効率的です。
  • 整数キー(Integer Keys): ロックは1つまたは2つのBIGINT値によって識別されます。

実装

import { Pool, PoolClient } from 'pg';

interface AdvisoryLockOptions {
  lockId: number; // A unique BIGINT identifier for the lock
  timeoutMs?: number; // How long to wait for the lock (0 for non-blocking)
}

class PgAdvisoryLock {
  private pool: Pool;

  constructor(connectionString: string) {
    this.pool = new Pool({ connectionString });
  }

  /**
   * Acquires a transaction-level advisory lock.
   * The lock is automatically released when the transaction commits or rolls back.
   * @returns The PoolClient if the lock is acquired, otherwise null.
   */
  public async acquireTransactionLock(options: AdvisoryLockOptions): Promise<PoolClient | null> {
    const { lockId, timeoutMs = 0 } = options;
    const client = await this.pool.connect();

    try {
      await client.query('BEGIN'); // Start a transaction

      let lockAcquired = false;
      if (timeoutMs === 0) {
        // Non-blocking attempt
        const res = await client.query('SELECT pg_try_advisory_xact_lock($1)', [lockId]);
        lockAcquired = res.rows[0].pg_try_advisory_xact_lock;
      } else {
        // Blocking attempt with timeout
        // pg_advisory_xact_lock will block until acquired or connection reset.
        // We simulate a timeout by running it in a separate promise and racing.
        const acquirePromise = client.query('SELECT pg_advisory_xact_lock($1)', [lockId]);
        const timeoutPromise = new Promise<void>(resolve => setTimeout(() => resolve(), timeoutMs));

        await Promise.race([acquirePromise, timeoutPromise]);
        // If acquirePromise resolved, lockAcquired is true. If timeoutPromise resolved first, it's false.
        // This is a simplification; a more robust solution might involve a separate connection for the timeout check.
        // For pg_advisory_xact_lock, it either succeeds or blocks indefinitely.
        // pg_try_advisory_xact_lock is generally preferred for non-blocking or explicit retry logic.
        // For a true blocking with timeout, one might need a loop with pg_try_advisory_xact_lock.
        const res = await client.query('SELECT pg_try_advisory_xact_lock($1)', [lockId]); // Re-check after potential block
        lockAcquired = res.rows[0].pg_try_advisory_xact_lock;
      }

      if (lockAcquired) {
        console.log(`Transaction lock ${lockId} acquired.`);
        return client; // Return the client to the caller for transaction operations
      } else {
        await client.query('ROLLBACK'); // Rollback if lock not acquired
        client.release();
        console.warn(`Failed to acquire transaction lock ${lockId}.`);
        return null;
      }
    } catch (error) {
      await client.query('ROLLBACK');
      client.release();
      console.error(`Error acquiring transaction lock ${lockId}: ${error}`);
      throw error;
    }
  }

  /**
   * Releases a transaction-level advisory lock by committing or rolling back the transaction.
   * This method should be called by the client that acquired the lock.
   */
  public async releaseTransactionLock(client: PoolClient, success: boolean): Promise<void> {
    try {
      if (success) {
        await client.query('COMMIT');
        console.log('Transaction committed, lock released.');
      } else {
        await client.query('ROLLBACK');
        console.log('Transaction rolled back, lock released.');
      }
    } catch (error) {
      console.error(`Error releasing transaction lock: ${error}`);
      await client.query('ROLLBACK'); // Ensure rollback on error
    } finally {
      client.release();
    }
  }

  /**
   * Acquires a session-level advisory lock.
   * The lock persists until explicitly released or the session ends.
   * @returns true if lock is acquired, false otherwise.
   */
  public async acquireSessionLock(options: AdvisoryLockOptions): Promise<boolean> {
    const { lockId, timeoutMs = 0 } = options;
    const client = await this.pool.connect(); // New client for session lock

    try {
      let lockAcquired = false;
      if (timeoutMs === 0) {
        const res = await client.query('SELECT pg_try_advisory_lock($1)', [lockId]);
        lockAcquired = res.rows[0].pg_try_advisory_lock;
      } else {
        // For session locks, pg_advisory_lock blocks.
        // A timeout would require a separate mechanism or loop with pg_try_advisory_lock.
        // For simplicity, we'll use pg_try_advisory_lock with retries for timeout simulation.
        const startTime = Date.now();
        while (Date.now() - startTime < timeoutMs) {
          const res = await client.query('SELECT pg_try_advisory_lock($1)', [lockId]);
          if (res.rows[0].pg_try_advisory_lock) {
            lockAcquired = true;
            break;
          }
          await new Promise(resolve => setTimeout(resolve, 50)); // Small delay before retry
        }
      }

      if (lockAcquired) {
        console.log(`Session lock ${lockId} acquired.`);
        // Store the client if you need to explicitly release it later
        // For this example, we'll just return true and assume the caller manages the client.
        // In a real app, you'd likely keep a map of lockId -> client.
        return true;
      } else {
        client.release(); // Release client if lock not acquired
        console.warn(`Failed to acquire session lock ${lockId}.`);
        return false;
      }
    } catch (error) {
      client.release();
      console.error(`Error acquiring session lock ${lockId}: ${error}`);
      throw error;
    }
  }

  /**
   * Releases a session-level advisory lock.
   * Requires the same client that acquired the lock.
   */
  public async releaseSessionLock(lockId: number, client: PoolClient): Promise<void> {
    try {
      await client.query('SELECT pg_advisory_unlock($1)', [lockId]);
      console.log(`Session lock ${lockId} released.`);
    } catch (error) {
      console.error(`Error releasing session lock ${lockId}: ${error}`);
    } finally {
      client.release();
    }
  }

  public async disconnect(): Promise<void> {
    await this.pool.end();
  }
}

// Example Usage (requires running PostgreSQL)
async function runPgAdvisoryLockExample() {
  const connectionString = 'postgresql://user:password@localhost:5432/mydatabase';
  const pgLock = new PgAdvisoryLock(connectionString);
  const lockId = 12345;

  // --- Transaction-level lock example ---
  console.log('\n--- Transaction-level lock ---');
  let client1: PoolClient | null = null;
  try {
    client1 = await pgLock.acquireTransactionLock({ lockId, timeoutMs: 100 });
    if (client1) {
      console.log('Client 1 acquired transaction lock. Performing database operations...');
      // Simulate database operation within the transaction
      await client1.query('INSERT INTO some_table (data) VALUES ($1)', ['transactional_data_1']);
      await new Promise(resolve => setTimeout(resolve, 1000)); // Simulate work
      await pgLock.releaseTransactionLock(client1, true); // Commit and release
    } else {
      console.log('Client 1 failed to acquire transaction lock.');
    }
  } catch (error) {
    console.error('Transaction example error:', error);
    if (client1) await pgLock.releaseTransactionLock(client1, false); // Rollback
  }

  // --- Session-level lock example ---
  console.log('\n--- Session-level lock ---');
  let sessionClient: PoolClient | null = null;
  try {
    sessionClient = await pgLock.pool.connect(); // Get a client for the session lock
    const acquired = await pgLock.acquireSessionLock({ lockId: lockId + 1, timeoutMs: 100 }); // Use a different lock ID

    if (acquired) {
      console.log('Client 1 acquired session lock. Performing operations...');
      // Simulate work
      await new Promise(resolve => setTimeout(resolve, 1500));

      // Attempt to acquire the same session lock from another client (should fail)
      console.log('Client 2 attempting to acquire the same session lock...');
      const client2 = await pgLock.pool.connect();
      const acquired2 = await pgLock.acquireSessionLock({ lockId: lockId + 1, timeoutMs: 100 });
      if (!acquired2) {
        console.log('Client 2 correctly failed to acquire the session lock.');
      }
      client2.release(); // Release client2 connection

      await pgLock.releaseSessionLock(lockId + 1, sessionClient);
    } else {
      console.log('Client 1 failed to acquire session lock.');
      if (sessionClient) sessionClient.release();
    }
  } catch (error) {
    console.error('Session example error:', error);
    if (sessionClient) sessionClient.release();
  }

  await pgLock.disconnect();
}

// To run this, you'd need a PostgreSQL instance running, e.g.:
// docker run -p 5432:5432 --name pg-lock -e POSTGRES_USER=user -e POSTGRES_PASSWORD=password -e POSTGRES_DB=mydatabase -d postgres
// Then: node your_script.js
// runPgAdvisoryLockExample();

分散協調のためのetcdリース

etcdは、高可用性と一貫性のある設定データストレージおよび協調サービスのために設計された分散キーバリューストアです。その核となる強みは、Raft合意アルゴリズムを使用していることで、強力な一貫性と耐障害性を提供します。etcdリースは、分散ロックのための強力なプリミティブです。

etcdリースの概要

etcdリースは、キーに関連付けられた有効期限(TTL)メカニズムです。キーがリースにアタッチされると、リースが期限切れになるとキーも期限切れになり、自動的に削除されます。クライアントは定期的に「キープアライブ」リクエストを送信することでリースを維持できます。これにより、堅牢な活性保証が提供されます。クライアントがクラッシュした場合、そのリースは期限切れになり、ロックは解放されます。

etcdリースによる分散ロックの実装

  1. リースを作成: クライアントは指定されたTTLでetcdにリースを要求します。
  2. ロックを取得: クライアントは、キーがすでに存在する場合は失敗するCREATE操作を使用して、そのリースに関連付けられた一時的なキー(例: /locks/my_resource)を作成しようとします。これはアトミックな「比較と交換」(CAS)操作です。
  3. キープアライブ: ロックが取得された場合、クライアントは定期的にそのリースに対するキープアライブリクエストを送信します。
  4. ロックを解放: クライアントは明示的にキーを削除するか、リースが期限切れになるのを待ちます。

etcdは、自然なフェンシングトークンとして機能するリビジョン番号も提供します。キーが作成されると、リビジョン番号が割り当てられます。その後の変更や削除は、このリビジョンをインクリメントします。クライアントはロックを取得し、そのリビジョンを取得し、このリビジョンを共有リソースに渡すことができます。リソースは、リビジョンをチェックすることで、操作が現在のロック保持者によって実行されていることを検証できます。

import { Etcd3, IOptions } from 'etcd3';

interface EtcdLockOptions {
  lockKey: string;
  ttlSeconds: number; // Lease time-to-live in seconds
  retryDelayMs?: number;
  maxRetries?: number;
}

class EtcdDistributedLock {
  private etcd: Etcd3;

  constructor(etcdHosts: string | string[], options?: IOptions) {
    this.etcd = new Etcd3({ hosts: etcdHosts, ...options });
  }

  /**
   * Acquires a distributed lock using etcd leases and compare-and-swap.
   * Returns the lease ID and the revision number (fencing token) if successful.
   */
  public async acquireLock(options: EtcdLockOptions): Promise<{ leaseId: string; revision: number } | null> {
    const { lockKey, ttlSeconds, retryDelayMs = 100, maxRetries = 10 } = options;
    const clientValue = Math.random().toString(36).substring(2, 15); // Unique client identifier

    let retries = 0;
    while (retries < maxRetries) {
      try {
        // 1. Create a lease
        const lease = this.etcd.lease(ttlSeconds);
        const leaseId = lease.id;

        // 2. Attempt to acquire the lock using a transaction (compare-and-swap)
        // Ensure the key does not exist, then create it with our value and lease.
        const transactionResult = await this.etcd.txn()
          .if(this.etcd.op(lockKey).not.exists())
          .then(this.etcd.op(lockKey).put(clientValue).withLease(leaseId))
          .commit();

        if (transactionResult.succeeded) {
          // Lock acquired. Start keep-alive.
          lease.on('lost', () => {
            console.error(`Lease ${leaseId} for lock ${lockKey} lost!`);
            // In a real application, this would trigger recovery or error handling.
          });
          lease.on('expired', () => {
            console.warn(`Lease ${leaseId} for lock ${lockKey} expired!`);
          });
          lease.on('keepalive established', () => {
            // console.log(`Keep-alive established for lease ${leaseId}`);
          });
          lease.on('keepalive error', (err) => {
            console.error(`Keep-alive error for lease ${leaseId}: ${err}`);
          });

          // Get the revision number for fencing
          const getResult = await this.etcd.get(lockKey).string();
          if (getResult) {
            const revision = getResult.modRevision;
            console.log(`Lock '${lockKey}' acquired by client '${clientValue}' with lease ${leaseId}, revision ${revision}`);
            return { leaseId, revision };
          } else {
            // This should ideally not happen if transaction succeeded, but defensive check.
            console.error(`Failed to get revision for lock ${lockKey} after acquiring.`);
            await this.releaseLock(lockKey, leaseId); // Clean up
            return null;
          }
        } else {
          // Lock already held by another client
          console.warn(`Lock '${lockKey}' already held. Retrying...`);
          await lease.revoke(); // Revoke our unused lease
          retries++;
          await new Promise(resolve => setTimeout(resolve, retryDelayMs));
        }
      } catch (error) {
        console.error(`Error acquiring lock '${lockKey}': ${error}`);
        retries++;
        await new Promise(resolve => setTimeout(resolve, retryDelayMs));
      }
    }

    console.error(`Failed to acquire lock '${lockKey}' after ${maxRetries} retries.`);
    return null;
  }

  /**
   * Releases a distributed lock by revoking its lease.
   * This will automatically delete the associated key.
   */
  public async releaseLock(lockKey: string, leaseId: string): Promise<void> {
    try {
      // Revoking the lease automatically deletes all keys associated with it.
      await this.etcd.lease.revoke(leaseId);
      console.log(`Lock '${lockKey}' with lease ${leaseId} released.`);
    } catch (error) {
      console.error(`Error releasing lock '${lockKey}' with lease ${leaseId}: ${error}`);
    }
  }

  /**
   * Verifies the fencing token (revision) for an operation.
   * The resource being protected should call this before performing an operation.
   */
  public async verifyFencingToken(lockKey: string, expectedRevision: number): Promise<boolean> {
    try {
      const getResult = await this.etcd.get(lockKey).string();
      if (!getResult) {
        console.warn(`Fencing check failed: Lock key '${lockKey}' does not exist.`);
        return false;
      }
      if (getResult.modRevision !== expectedRevision) {
        console.warn(`Fencing check failed: Expected revision ${expectedRevision}, got ${getResult.modRevision}.`);
        return false;
      }
      return true;
    } catch (error) {
      console.error(`Error during fencing token verification for '${lockKey}': ${error}`);
      return false;
    }
  }

  public async disconnect(): Promise<void> {
    await this.etcd.close();
  }
}

// Example Usage (requires running etcd)
async function runEtcdLockExample() {
  const etcdHosts = ['localhost:2379']; // Or an array of hosts for a cluster
  const etcdLock = new EtcdDistributedLock(etcdHosts);

  const lockKey = '/app/resource/processor_lock';
  const ttlSeconds = 5; // Lease expires in 5 seconds if not kept alive

  console.log('Attempting to acquire lock 1...');
  const lockInfo1 = await etcdLock.acquireLock({ lockKey, ttlSeconds });

  if (lockInfo1) {
    console.log(`Client 1 acquired lock with lease ${lockInfo1.leaseId}, revision ${lockInfo1.revision}.`);
    // Simulate critical operation with fencing check
    const isFenced = await etcdLock.verifyFencingToken(lockKey, lockInfo1.revision);
    if (isFenced) {
      console.log('Fencing token verified. Performing critical operation...');
      await new Promise(resolve => setTimeout(resolve, 3000)); // Simulate work
      console.log('Client 1 finished critical operation. Releasing lock...');
      await etcdLock.releaseLock(lockKey, lockInfo1.leaseId);
    } else {
      console.error('Fencing check failed for client 1. Aborting operation.');
      await etcdLock.releaseLock(lockKey, lockInfo1.leaseId); // Ensure release
    }
  } else {
    console.log('Client 1 failed to acquire lock.');
  }

  console.log('\nAttempting to acquire lock 2 (should succeed after lock 1 is released)...');
  const lockInfo2 = await etcdLock.acquireLock({ lockKey, ttlSeconds });

  if (lockInfo2) {
    console.log(`Client 2 acquired lock with lease ${lockInfo2.leaseId}, revision ${lockInfo2.revision}.`);
    await new Promise(resolve => setTimeout(resolve, 1000));
    await etcdLock.releaseLock(lockKey, lockInfo2.leaseId);
  } else {
    console.log('Client 2 failed to acquire lock.');
  }

  await etcdLock.disconnect();
}

// To run this, you'd need an etcd instance running, e.g.:
// docker run -p 2379:2379 -p 2380:2380 --name etcd-single -d quay.io/coreos/etcd:latest etcd -advertise-client-urls http://0.0.0.0:2379 -listen-client-urls http://0.0.0.0:2379
// Then: node your_script.js
// runEtcdLockExample();
Advertisement

アーキテクチャ比較

機能Redlock (Redis)PostgreSQLアドバイザリロックetcdリース
一貫性モデル結果整合性 (複数インスタンスRedis)強力な一貫性 (PostgreSQLのACID保証)強力な一貫性 (Raft合意)
耐障害性Redisインスタンスの過半数クォーラムPostgreSQLクラスタ (例: ストリーミングレプリケーション)etcdノードの過半数クォーラム
フェンシングサポートネイティブではない; 外部メカニズムが必要ネイティブではない; アプリケーションロジックでシミュレート可能リビジョン番号によるネイティブサポート
活性保証TTLベースの有効期限トランザクション/セッション終了、または明示的なロック解除キープアライブ付きリースTTL
オーバーヘッドRedisインスタンスへの複数ネットワークラウンドトリップ単一のデータベース接続/トランザクションetcdクラスタへのネットワークラウンドトリップ
複雑性中程度 (クライアント側のクォーラムロジック)低 (SQL関数)中程度 (etcdクライアント、リース管理)
ユースケース高スループット、低レイテンシ、非クリティカルなロックデータベース中心のアプリケーション、トランザクションに紐づく作業クリティカルな協調、リーダー選出、設定管理
依存関係RedisクラスタPostgreSQLデータベースetcdクラスタ

本番環境での落とし穴とトラブルシューティング

  1. クロック同期 (Redlock):

    • 落とし穴: Redisインスタンス間、またはクライアントとRedis間の大幅なクロックスキューは、ロックの早期期限切れや意図よりも長くロックが保持されることにつながり、相互排他を侵害します。
    • 解決策: すべてのノードで厳密なクロック同期のためにNTPまたはPTPを実装します。Redlockのvalidity_timeはこれを緩和しようとしますが、完全な解決策ではありません。高信頼性システムでは、Redlockの使用を避けるべきです。
  2. GCポーズ/プロセス停止 (すべて):

    • 落とし穴: ロックを保持しているクライアントが長いGCポーズを経験します。そのロックは期限切れになり、別のクライアントがそれを取得し、最初のクライアントが再開して、古い状態に基づいて操作したり、新しいロック保持者と競合したりします。
    • 解決策: フェンシングトークンを実装します。etcdの場合はリビジョン番号を使用します。PostgreSQLの場合は、専用のロックテーブルにversion列を使用し、ロック取得ごとにインクリメントすることができます。Redlockの場合、これは外部協調なしでは根本的な弱点です。さらに、アプリケーションの一時停止時間を監視し、JVM/ランタイム設定を調整します。
  3. ネットワークパーティション (すべて):

    • 落とし穴: クライアントがロックを取得した後、ロックサービスの大半からパーティション化されます。クライアントはまだロックを保持していると信じますが、健全なパーティションでは新しいロックが付与されます。
    • 解決策: これは、合意アルゴリズム(etcdのRaft)やクォーラムベースのアプローチ(Redlock)が役立つように設計されている部分です。しかし、フェンシングは依然として重要です。パーティション化されたクライアントは、ロックを持っていると「思っていても」、保護しようとしているリソースによって行動を阻止されなければなりません。
  4. リース/TTL管理 (Redlock, etcd):

    • 落とし穴: TTLの誤った設定やキープアライブの送信失敗は、ロックの早期期限切れや無期限の保持につながる可能性があります。
    • 解決策: ネットワークレイテンシと予想される操作期間を考慮して、TTLを慎重に選択します。エラー処理とリトライロジックを備えた堅牢なキープアライブループを実装します。リースがアクティブな場合でも、完了したらクライアントが明示的にロックを解放するようにします。
  5. PostgreSQLトランザクション分離:

    • 落とし穴: pg_advisory_lock(セッションレベル)の代わりにpg_advisory_xact_lock(トランザクションレベル)を使用すると、ロックがトランザクションのスコープを超えて持続し、明示的に解放されない場合にデッドロックや予期しない動作を引き起こす可能性があります。
    • 解決策: 自然にトランザクション的な操作には常にpg_advisory_xact_lockを優先します。セッションレベルロックを使用する場合は、finallyブロックまたは同様のクリーンアップメカニズムで明示的なpg_advisory_unlock呼び出しを確実に実行します。
  6. リソース枯渇 (PostgreSQL):

    • 落とし穴: アドバイザリロックの過度な使用は、特にクライアントがロックを待ってブロックする場合、接続プールリソースを消費する可能性があります。
    • 解決策: 非ブロッキング試行にはpg_try_advisory_xact_lockを使用し、バックオフ付きのクライアント側リトライロジックを実装します。接続プールの使用状況を監視し、max_connectionsを調整します。

よくある質問

Q1: RedlockとetcdまたはPostgreSQLアドバイザリロックは、いつ使い分けるべきですか?

A1: Redlockは、低レイテンシ、高スループットのロックが強く求められ、その既知の安全性上の制限(特にフェンシングとクロックドリフトに関して)を受け入れるか、複雑な外部フェンシングメカニズムを実装する意思がある場合にのみ使用してください。これは、相互排他の偶発的な違反が許容される、非クリティカルな「ベストエフォート」ロックに最適です。強力な一貫性と安全性、特にフェンシングが必要な場合は、etcdまたはPostgreSQLアドバイザリロックが優れています。共有リソースが主にデータベース自体であり、トランザクションに紐づくロックが必要な場合はPostgreSQLを選択してください。専用の強力に一貫性のあるキーバリューストアが適切な、より広範な分散協調、リーダー選出、設定管理にはetcdを選択してください。

Q2: フェンシングトークンは、古いクライアントが害を及ぼすのをどのように本当に防ぐのですか?

A2: フェンシングトークンは、保護されているリソースがトークンを検証することを要求することで機能します。クライアントがロックを取得すると、一意で単調増加するトークン(例:etcdのリビジョン番号)を受け取ります。クライアントが共有リソースを変更しようとするとき、このトークンを提示しなければなりません。リソース(例:ストレージサービス、メッセージキュー)は、提示されたトークンがそのリソースにとって最新の有効なトークンであるかどうかをチェックします。もし古いトークンが提示された場合(つまり、別のクライアントがより高いトークンでロックを取得している場合)、操作は拒否されます。これにより、古いクライアントが誤ってロックを保持していると信じていても、データを破損するのを防ぎます。

Q3: 単一のRedisインスタンスを分散ロックに使用できますか?

A3: いいえ、いかなるレベルの耐障害性も要求される本番システムでは絶対にできません。単一のRedisインスタンスは単一障害点です。それがクラッシュすると、すべてのロックが失われ、クライアントが同時に処理を進め、データ破損につながる可能性があります。Redlockはインスタンスのクォーラムでこれを緩和しようとしますが、それでも制限があります。真剣な分散ロックには、耐障害性のあるバックエンド(Redlockを備えたRedisクラスタ、PostgreSQLクラスタ、またはetcdクラスタなど)が必須です。

Q4: 分散ロックのパフォーマンスへの影響は何ですか?

A4: 分散ロックは、ネットワークラウンドトリップと協調のオーバーヘッドにより、本質的にレイテンシを導入します。

  • Redlock: Redisインスタンスが近くにある場合、取得は高速ですが、複数のネットワーク呼び出しを伴います。
  • PostgreSQLアドバイザリロック: データベース接続がすでに確立されており、データベースが重い負荷下にない場合、非常に高速です。トランザクションレベルロックは、既存のトランザクションに紐づいているため特に効率的です。
  • etcdリース: etcdクラスタへのネットワーク呼び出しを伴い、Raft合意を使用するため、単一のRedisインスタンスと比較して多少のレイテンシが追加されます。しかし、etcdはこのために最適化されており、強力な保証を提供します。

パフォーマンスへの影響は、ロックの取得/解放の頻度、ネットワークレイテンシ、および基盤となるロックサービスの負荷に大きく依存します。分散ロックによって保護されるクリティカルセクションを最小限に抑えるようにシステムを設計してください。

Q5: クライアントがロックを保持している最中にクラッシュした場合、どうなりますか?

A5: これは分散ロックの主要な懸念事項であり、その活性保証によって対処されます。

  • Redlock: ロックは最終的にTTLによって期限切れになります。random_valueは、クライアントが再起動してロックを解放しようとした場合、自身のロックのみを解放し、別のクライアントが取得した新しいロックは解放しないことを保証します。
  • PostgreSQLアドバイザリロック: トランザクションレベルロック(pg_advisory_xact_lock)の場合、クライアントのデータベース接続が終了すると(例:クラッシュによる)、トランザクションがロールバックされ、ロックは自動的に解放されます。セッションレベルロック(pg_advisory_lock)の場合、セッションが終了すると解放されます。
  • etcdリース: クライアントがクラッシュしてキープアライブの送信を停止すると、ロックキーに関連付けられたリースが期限切れになります。リースが期限切れになると、etcdは自動的にロックキーを削除し、ロックを他のクライアントが利用できるようにします。これは非常に堅牢なメカニズムです。

すべての場合において、ロックは最終的に解放され、デッドロックを防ぎます。しかし、クラッシュしたクライアントが回復し、古いロックで操作しようとした場合に害を及ぼすのを防ぐためには、フェンシングが依然として必要です。

Share this article:

Stay Updated

Get the latest posts delivered straight to your inbox.

Free Developer Utilities

Free In-Browser Developer Tools

Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.

Explore Tools
Advertisement