@ThreadSafe public class CoordinatorStateDAO extends Object
CoordinatorState operations to the appropriate
DDB table based on the current TableMigrationStatus.
leaseTableDaoDelegate field
for its reads/writes and sets the TableMigrationStatusProvider before calling initialize().initialize() to have been called first.TableMigrationStateMachine.handleLeaderLockResult
| Constructor and Description |
|---|
CoordinatorStateDAO(software.amazon.awssdk.services.dynamodb.DynamoDbAsyncClient dynamoDbAsyncClient,
CoordinatorConfig.CoordinatorStateTableConfig coordinatorStateTableConfig,
String leaseTableName,
TableMigrationStatusProvider tableMigrationStatusProvider) |
| Modifier and Type | Method and Description |
|---|---|
boolean |
createCoordinatorStateIfNotExists(@NonNull CoordinatorState state)
Create a coordinator state if it does not already exist.
|
boolean |
deleteCoordinatorState(@NonNull String key)
Delete a coordinator state by key.
|
void |
executeTransactWrite(List<software.amazon.awssdk.services.dynamodb.model.TransactWriteItem> transactWriteItems)
Execute a DynamoDB TransactWriteItems request with the given list of
TransactWriteItems. |
CoordinatorState |
getCoordinatorState(@NonNull String key)
Get a single
CoordinatorState by key. |
com.amazonaws.services.dynamodbv2.AmazonDynamoDBLockClient |
getDDBLockClient()
Get the active DDB lock client based on the current table migration status.
|
LeaseTableCoordinatorStateDAODelegate |
getLeaseTableDaoDelegate() |
LegacyTableCoordinatorStateDAODelegate |
getLegacyTableDaoDelegate() |
void |
initialize()
Initialize the DAO for write operations.
|
void |
initializeDelegates()
Initialize the delegates for read operations.
|
void |
initializeLockClients(long leaseDurationMillis,
long heartbeatPeriodMillis,
String workerId)
Initialize the DDB lock clients for leader election.
|
List<CoordinatorState> |
listCoordinatorState()
List all coordinator states.
|
List<CoordinatorState> |
listCoordinatorStateByEntityType(EntityType.CoordinatorStateType entityType)
List coordinator states by entity type.
|
void |
shutdownLockClients()
Shutdown both lock clients, closing their background heartbeat threads.
|
boolean |
updateCoordinatorStateWithExpectation(@NonNull CoordinatorState state,
Map<String,software.amazon.awssdk.services.dynamodb.model.ExpectedAttributeValue> expectations)
Update a coordinator state with expectations.
|
public CoordinatorStateDAO(software.amazon.awssdk.services.dynamodb.DynamoDbAsyncClient dynamoDbAsyncClient,
CoordinatorConfig.CoordinatorStateTableConfig coordinatorStateTableConfig,
String leaseTableName,
TableMigrationStatusProvider tableMigrationStatusProvider)
public LegacyTableCoordinatorStateDAODelegate getLegacyTableDaoDelegate()
public LeaseTableCoordinatorStateDAODelegate getLeaseTableDaoDelegate()
public void initializeDelegates()
throws DependencyException
DependencyException - if unable to determine legacy table existencepublic void initialize()
throws InvalidStateException
initializeDelegates() and after the TableMigrationStatusProvider
has moved past UNKNOWN.InvalidStateException - if the TableMigrationStatusProvider is still UNKNOWNpublic CoordinatorState getCoordinatorState(@NonNull @NonNull String key) throws DependencyException, InvalidStateException, ProvisionedThroughputException
CoordinatorState by key.
Uses the fallback read pattern: if status != COMPLETE, tries legacy first then lease table.
If status == COMPLETE, reads only from lease table.key - the coordinator state keyDependencyException - if DDB fails unexpectedlyInvalidStateException - if a required table does not existProvisionedThroughputException - if DDB lacks capacitypublic List<CoordinatorState> listCoordinatorState() throws ProvisionedThroughputException, DependencyException, InvalidStateException
DependencyException - if DDB fails unexpectedlyInvalidStateException - if a required table does not existProvisionedThroughputException - if DDB lacks capacitypublic List<CoordinatorState> listCoordinatorStateByEntityType(EntityType.CoordinatorStateType entityType) throws ProvisionedThroughputException, DependencyException, InvalidStateException
entityType - the entity type to filter byDependencyException - if DDB fails unexpectedlyInvalidStateException - if a required table does not existProvisionedThroughputException - if DDB lacks capacitypublic boolean createCoordinatorStateIfNotExists(@NonNull
@NonNull CoordinatorState state)
throws DependencyException,
InvalidStateException,
ProvisionedThroughputException
state - the state to createDependencyException - if DDB fails unexpectedlyInvalidStateException - if not initialized or table does not existProvisionedThroughputException - if DDB lacks capacitypublic boolean updateCoordinatorStateWithExpectation(@NonNull
@NonNull CoordinatorState state,
Map<String,software.amazon.awssdk.services.dynamodb.model.ExpectedAttributeValue> expectations)
throws DependencyException,
InvalidStateException,
ProvisionedThroughputException
state - the state to updateexpectations - conditional expectations for the updateDependencyException - if DDB fails unexpectedlyInvalidStateException - if not initialized or table does not existProvisionedThroughputException - if DDB lacks capacitypublic boolean deleteCoordinatorState(@NonNull
@NonNull String key)
throws ProvisionedThroughputException,
InvalidStateException,
DependencyException
key - the key to deleteDependencyException - if DDB fails unexpectedlyInvalidStateException - if not initialized or table does not existProvisionedThroughputException - if DDB lacks capacitypublic void initializeLockClients(long leaseDurationMillis,
long heartbeatPeriodMillis,
String workerId)
The lock clients have background heartbeat threads that keep acquired locks alive. An idle lock client (one that hasn't acquired a lock) will have a running heartbeat thread but it will be effectively a no-op since there are no locks to heartbeat.
This method is called once during startup and is not thread-safe with respect to
concurrent calls. It must be called before any calls to getDDBLockClient().
leaseDurationMillis - the lock lease duration in millisecondsheartbeatPeriodMillis - the heartbeat period in millisecondsworkerId - the owner name for the lockpublic com.amazonaws.services.dynamodbv2.AmazonDynamoDBLockClient getDDBLockClient()
Thread-safe: the fields are assigned once during initializeLockClients(long, long, java.lang.String) (startup)
and never reassigned. The routing decision is based on the TableMigrationStatusProvider
which is updated atomically. The caller (isLeader()) is synchronized,
providing happens-before guarantees.
AmazonDynamoDBLockClient for the current migration statepublic void shutdownLockClients()
public void executeTransactWrite(List<software.amazon.awssdk.services.dynamodb.model.TransactWriteItem> transactWriteItems) throws DependencyException
TransactWriteItems.
Used by the table migration state machine to atomically move entries between tables
(put into lease table + delete from legacy table in a single transaction).transactWriteItems - the list of transact write items (max 100 per DDB limit)DependencyException - if DDB fails unexpectedly or the transaction is cancelledCopyright © 2026. All rights reserved.