This is an automated email from the ASF dual-hosted git repository.
jerryshao pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/main by this push:
new 26a9fff5e6 [#13387] refactor(core): add metadata environment profile
(#13367)
26a9fff5e6 is described below
commit 26a9fff5e64d790cabe5570a8ce1c863fb6a38a2
Author: mchades <[email protected]>
AuthorDate: Tue Sep 22 11:27:08 2026 +0800
[#13387] refactor(core): add metadata environment profile (#13367)
### What changes were proposed in this pull request?
- Split `GravitinoEnv` initialization into common components, internal
metadata components, and server-side integrations.
- Add `initializeMetadataComponents(Config)` as the entry point for
initializing the internal metadata layer.
- Refactor full-server initialization to reuse the same metadata
initialization path before adding hooks, events, auxiliary services, and
jobs.
- Add a normalized internal partition dispatcher.
- Add tests covering component composition, lifecycle behavior,
server-integration isolation, and authorization modes.
### Why are the changes needed?
`GravitinoEnv` currently initializes internal metadata components and
server-side integrations in the same method. This mixes responsibilities
and makes component dependencies and lifecycle ownership difficult to
understand and maintain.
This change makes the initialization hierarchy explicit:
1. Common infrastructure.
2. Internal metadata components.
3. Server-side integrations.
The full-server profile continues to initialize the same components, but
now composes them through clearly defined initialization layers. This
reduces coupling, avoids duplicated initialization logic, and provides a
stable metadata initialization boundary for future reuse.
Fix: #13387
### Does this PR introduce _any_ user-facing change?
No REST API or configuration changes.
`GravitinoEnv` adds `initializeMetadataComponents(Config)` and
`internalPartitionDispatcher()` for internal or embedded use. Existing
base and full initialization behavior remains unchanged.
### How was this patch tested?
- `./gradlew :core:spotlessApply`
- `./gradlew :core:check :core:javadoc -PskipITs`
- `./gradlew compileDistribution -PskipITs -x test`
- `git diff --check`
---
.../java/org/apache/gravitino/GravitinoEnv.java | 398 ++++++++++++++-------
.../TestGravitinoEnvMetadataComponents.java | 336 +++++++++++++++++
2 files changed, 602 insertions(+), 132 deletions(-)
diff --git a/core/src/main/java/org/apache/gravitino/GravitinoEnv.java
b/core/src/main/java/org/apache/gravitino/GravitinoEnv.java
index 5b38b7182e..36764b65b5 100644
--- a/core/src/main/java/org/apache/gravitino/GravitinoEnv.java
+++ b/core/src/main/java/org/apache/gravitino/GravitinoEnv.java
@@ -144,6 +144,7 @@ public class GravitinoEnv {
private TableDispatcher internalTableDispatcher;
private PartitionDispatcher partitionDispatcher;
+ private PartitionDispatcher internalPartitionDispatcher;
private FilesetDispatcher filesetDispatcher;
@@ -206,6 +207,7 @@ public class GravitinoEnv {
private FutureGrantManager futureGrantManager;
private GravitinoAuthorizer gravitinoAuthorizer;
private StatisticDispatcher statisticDispatcher;
+ private StatisticDispatcher internalStatisticDispatcher;
protected GravitinoEnv() {}
@@ -229,14 +231,35 @@ public class GravitinoEnv {
*/
public void initializeBaseComponents(Config config) {
LOG.info("Initializing Gravitino base environment...");
- this.config = config;
- FileFetcher.get().initialize(config.get(Configs.BLOCK_UNSAFE_REMOTE_URI));
- SecretPropertyUtils.configureSensitiveKeyKeywords(config);
+ initializeConfig(config);
this.manageFullComponents = false;
initBaseComponents();
LOG.info("Gravitino base environment is initialized.");
}
+ /**
+ * Initializes components required for normalized metadata operations.
+ *
+ * <p>This initialization profile does not initialize event listeners, audit
logging, metadata
+ * hooks, auxiliary services, or job management.
+ *
+ * <p>This method must be called on {@link #getInstance()}. Some metadata
components read their
+ * dependencies directly from that singleton instead of from the object
being initialized.
+ *
+ * @param config The configuration object to initialize the environment.
+ */
+ public void initializeMetadataComponents(Config config) {
+ Preconditions.checkState(
+ this == getInstance(),
+ "Metadata components must be initialized on
GravitinoEnv.getInstance().");
+ LOG.info("Initializing Gravitino metadata environment...");
+ initializeConfig(config);
+ this.manageFullComponents = false;
+ initCommonComponents();
+ initMetadataComponents();
+ LOG.info("Gravitino metadata environment is initialized.");
+ }
+
/**
* Initialize all components, used for Gravitino server.
*
@@ -244,9 +267,7 @@ public class GravitinoEnv {
*/
public void initializeFullComponents(Config config) {
LOG.info("Initializing Gravitino full environment...");
- this.config = config;
- FileFetcher.get().initialize(config.get(Configs.BLOCK_UNSAFE_REMOTE_URI));
- SecretPropertyUtils.configureSensitiveKeyKeywords(config);
+ initializeConfig(config);
this.manageFullComponents = true;
initBaseComponents();
initGravitinoServerComponents();
@@ -427,6 +448,19 @@ public class GravitinoEnv {
return partitionDispatcher;
}
+ /**
+ * Get the internal PartitionDispatcher associated with the Gravitino
environment.
+ *
+ * <p>The internal dispatcher preserves normalization but skips event
emission.
+ *
+ * @return The internal PartitionDispatcher instance.
+ */
+ public PartitionDispatcher internalPartitionDispatcher() {
+ Preconditions.checkArgument(
+ internalPartitionDispatcher != null, "GravitinoEnv is not
initialized.");
+ return internalPartitionDispatcher;
+ }
+
/**
* Get the FilesetDispatcher associated with the Gravitino environment.
*
@@ -739,13 +773,26 @@ public class GravitinoEnv {
return statisticDispatcher;
}
+ /**
+ * Get the internal StatisticDispatcher associated with the Gravitino
environment.
+ *
+ * @return The internal StatisticDispatcher instance.
+ */
+ public StatisticDispatcher internalStatisticDispatcher() {
+ Preconditions.checkArgument(
+ internalStatisticDispatcher != null, "GravitinoEnv is not
initialized.");
+ return internalStatisticDispatcher;
+ }
+
public boolean cacheEnabled() {
return config == null || config.get(Configs.CACHE_ENABLED);
}
public void start() {
metricsSystem.start();
- eventListenerManager.start();
+ if (eventListenerManager != null) {
+ eventListenerManager.start();
+ }
if (manageFullComponents) {
auxServiceManager.serviceStart();
}
@@ -796,9 +843,11 @@ public class GravitinoEnv {
}
}
- if (statisticDispatcher != null) {
+ StatisticDispatcher statisticDispatcherToClose =
+ statisticDispatcher != null ? statisticDispatcher :
internalStatisticDispatcher;
+ if (statisticDispatcherToClose != null) {
try {
- statisticDispatcher.close();
+ statisticDispatcherToClose.close();
} catch (Exception e) {
LOG.warn("Failed to close StatisticDispatcher", e);
}
@@ -816,11 +865,7 @@ public class GravitinoEnv {
}
private void initBaseComponents() {
- this.kmsClientRegistry = new KmsClientRegistry(config);
- this.secretManager = new SecretManager(config);
-
- this.metricsSystem = new MetricsSystem();
- metricsSystem.register(new JVMMetricsSource());
+ initCommonComponents();
this.eventListenerManager = new EventListenerManager();
eventListenerManager.init(
@@ -831,77 +876,77 @@ public class GravitinoEnv {
auditLogManager.init(config, eventListenerManager);
}
- private void initGravitinoServerComponents() {
- // Initialize EntityStore
- this.entityStore = EntityStoreFactory.createEntityStore(config);
- entityStore.initialize(config);
+ private void initializeConfig(Config config) {
+ this.config = config;
+ FileFetcher.get().initialize(config.get(Configs.BLOCK_UNSAFE_REMOTE_URI));
+ SecretPropertyUtils.configureSensitiveKeyKeywords(config);
+ }
- // create and initialize a random id generator
- this.idGenerator = new RandomIdGenerator();
+ private void initCommonComponents() {
+ this.kmsClientRegistry = new KmsClientRegistry(config);
+ this.secretManager = new SecretManager(config);
- // Tree lock
- this.lockManager = new LockManager(config);
+ this.metricsSystem = new MetricsSystem();
+ metricsSystem.register(new JVMMetricsSource());
+ }
- // Create and initialize Catalog related modules first so MetalakeManager
can force-drop
- // child catalogs through CatalogManager.dropCatalog (same path as
FilesetCatalogOperations).
- // CatalogEventDispatcher -> CatalogNormalizeDispatcher ->
CatalogHookDispatcher ->
- // CatalogManager
- // CatalogManager registers its own change-log listener with the entity
store (when the store
- // supports it), so no poller wiring is needed here.
- this.catalogManager = new CatalogManager(config, entityStore, idGenerator,
secretManager);
+ private MetadataOperations initMetadataComponents() {
+ initEntityStoreAndCatalogManager();
- // Create and initialize metalake related modules, the operation chain is:
- // MetalakeEventDispatcher -> MetalakeNormalizeDispatcher ->
MetalakeHookDispatcher ->
- // MetalakeManager
this.metalakeManager = new MetalakeManager(entityStore, idGenerator,
catalogManager);
this.internalMetalakeDispatcher = new
MetalakeNormalizeDispatcher(metalakeManager);
- MetalakeHookDispatcher metalakeHookDispatcher = new
MetalakeHookDispatcher(metalakeManager);
- MetalakeNormalizeDispatcher metalakeNormalizeDispatcher =
- new MetalakeNormalizeDispatcher(metalakeHookDispatcher);
- this.metalakeDispatcher = new MetalakeEventDispatcher(eventBus,
metalakeNormalizeDispatcher);
-
this.internalCatalogDispatcher = new
CatalogNormalizeDispatcher(catalogManager);
- CatalogHookDispatcher catalogHookDispatcher = new
CatalogHookDispatcher(catalogManager);
- CatalogNormalizeDispatcher catalogNormalizeDispatcher =
- new CatalogNormalizeDispatcher(catalogHookDispatcher);
- this.catalogDispatcher = new CatalogEventDispatcher(eventBus,
catalogNormalizeDispatcher);
this.credentialOperationDispatcher =
new CredentialOperationDispatcher(catalogManager, entityStore,
idGenerator, secretManager);
-
this.secretPropertyOperationDispatcher =
new SecretPropertyOperationDispatcher(
catalogManager, entityStore, idGenerator, secretManager);
// Fileset dispatcher is created before schema dispatcher so schema can
take it directly.
+ FilesetOperationDispatcher filesetOperationDispatcher =
initInternalFilesetDispatcher();
+ SchemaOperationDispatcher schemaOperationDispatcher =
initInternalSchemaDispatcher();
+ initInternalTableDispatcher();
+ initInternalPartitionDispatcher();
+ TopicOperationDispatcher topicOperationDispatcher =
initInternalTopicDispatcher();
+ ModelOperationDispatcher modelOperationDispatcher =
initInternalModelDispatcher();
+ FunctionOperationDispatcher functionOperationDispatcher =
+ initInternalFunctionDispatcher(schemaOperationDispatcher);
+ initInternalViewDispatcher();
+ initSemanticModelDispatcher(schemaOperationDispatcher);
+
+ this.internalStatisticDispatcher = new StatisticManager(entityStore,
idGenerator, config);
+ initInternalAuthorizationComponents();
+
+ this.internalTagDispatcher = new TagManager(idGenerator, entityStore);
+ this.internalPolicyDispatcher = new PolicyManager(idGenerator,
entityStore);
+
+ return new MetadataOperations(
+ filesetOperationDispatcher,
+ schemaOperationDispatcher,
+ topicOperationDispatcher,
+ modelOperationDispatcher,
+ functionOperationDispatcher);
+ }
+
+ private FilesetOperationDispatcher initInternalFilesetDispatcher() {
FilesetOperationDispatcher filesetOperationDispatcher =
new FilesetOperationDispatcher(catalogManager, entityStore,
idGenerator, secretManager);
- FilesetNormalizeDispatcher internalFilesetNormalizeDispatcher =
+ this.internalFilesetDispatcher =
new FilesetNormalizeDispatcher(filesetOperationDispatcher,
catalogManager);
- this.internalFilesetDispatcher = internalFilesetNormalizeDispatcher;
- FilesetHookDispatcher filesetHookDispatcher =
- new FilesetHookDispatcher(filesetOperationDispatcher);
- FilesetNormalizeDispatcher filesetNormalizeDispatcher =
- new FilesetNormalizeDispatcher(filesetHookDispatcher, catalogManager);
- this.filesetDispatcher = new FilesetEventDispatcher(eventBus,
filesetNormalizeDispatcher);
+ return filesetOperationDispatcher;
+ }
+ private SchemaOperationDispatcher initInternalSchemaDispatcher() {
SchemaOperationDispatcher schemaOperationDispatcher =
new SchemaOperationDispatcher(
- catalogManager,
- entityStore,
- idGenerator,
- secretManager,
- internalFilesetNormalizeDispatcher);
+ catalogManager, entityStore, idGenerator, secretManager,
internalFilesetDispatcher);
this.internalSchemaDispatcher =
new SchemaNormalizeDispatcher(schemaOperationDispatcher,
catalogManager);
- SchemaHookDispatcher schemaHookDispatcher = new
SchemaHookDispatcher(schemaOperationDispatcher);
- SchemaNormalizeDispatcher schemaNormalizeDispatcher =
- new SchemaNormalizeDispatcher(schemaHookDispatcher, catalogManager);
- this.schemaDispatcher = new SchemaEventDispatcher(eventBus,
schemaNormalizeDispatcher);
+ return schemaOperationDispatcher;
+ }
- TableOperationDispatcher tableOperationDispatcher =
- new TableOperationDispatcher(catalogManager, entityStore, idGenerator,
secretManager);
- this.internalTableDispatcher = tableOperationDispatcher;
+ private void initInternalTableDispatcher() {
TableOperationDispatcher internalTableOperationDispatcher =
new TableOperationDispatcher(
catalogManager,
@@ -911,60 +956,42 @@ public class GravitinoEnv {
secretManager);
this.internalTableDispatcher =
new TableNormalizeDispatcher(internalTableOperationDispatcher,
catalogManager);
- TableHookDispatcher tableHookDispatcher =
- new TableHookDispatcher(tableOperationDispatcher,
this::internalOwnerDispatcher);
- TableNormalizeDispatcher tableNormalizeDispatcher =
- new TableNormalizeDispatcher(tableHookDispatcher, catalogManager);
- this.tableDispatcher = new TableEventDispatcher(eventBus,
tableNormalizeDispatcher);
+ }
- // TODO: We can install hooks when we need, we only supports ownership
post hook,
- // partition doesn't have ownership, so we don't need it now.
+ private void initInternalPartitionDispatcher() {
PartitionOperationDispatcher partitionOperationDispatcher =
new PartitionOperationDispatcher(catalogManager, entityStore,
idGenerator, secretManager);
- PartitionNormalizeDispatcher partitionNormalizeDispatcher =
+ this.internalPartitionDispatcher =
new PartitionNormalizeDispatcher(partitionOperationDispatcher,
catalogManager);
- this.partitionDispatcher = new PartitionEventDispatcher(eventBus,
partitionNormalizeDispatcher);
+ }
+ private TopicOperationDispatcher initInternalTopicDispatcher() {
TopicOperationDispatcher topicOperationDispatcher =
new TopicOperationDispatcher(catalogManager, entityStore, idGenerator,
secretManager);
- TopicNormalizeDispatcher internalTopicNormalizeDispatcher =
+ this.internalTopicDispatcher =
new TopicNormalizeDispatcher(topicOperationDispatcher, catalogManager);
- this.internalTopicDispatcher = internalTopicNormalizeDispatcher;
- TopicHookDispatcher topicHookDispatcher = new
TopicHookDispatcher(topicOperationDispatcher);
- TopicNormalizeDispatcher topicNormalizeDispatcher =
- new TopicNormalizeDispatcher(topicHookDispatcher, catalogManager);
- this.topicDispatcher = new TopicEventDispatcher(eventBus,
topicNormalizeDispatcher);
+ return topicOperationDispatcher;
+ }
+ private ModelOperationDispatcher initInternalModelDispatcher() {
ModelOperationDispatcher modelOperationDispatcher =
new ModelOperationDispatcher(catalogManager, entityStore, idGenerator,
secretManager);
- ModelNormalizeDispatcher internalModelNormalizeDispatcher =
+ this.internalModelDispatcher =
new ModelNormalizeDispatcher(modelOperationDispatcher, catalogManager);
- this.internalModelDispatcher = internalModelNormalizeDispatcher;
- ModelHookDispatcher modelHookDispatcher = new
ModelHookDispatcher(modelOperationDispatcher);
- ModelNormalizeDispatcher modelNormalizeDispatcher =
- new ModelNormalizeDispatcher(modelHookDispatcher, catalogManager);
- this.modelDispatcher = new ModelEventDispatcher(eventBus,
modelNormalizeDispatcher);
+ return modelOperationDispatcher;
+ }
- // Create and initialize Function related modules, the operation chain is:
- // FunctionEventDispatcher -> FunctionNormalizeDispatcher ->
FunctionHookDispatcher ->
- // FunctionOperationDispatcher
+ private FunctionOperationDispatcher initInternalFunctionDispatcher(
+ SchemaOperationDispatcher schemaOperationDispatcher) {
FunctionOperationDispatcher functionOperationDispatcher =
new FunctionOperationDispatcher(
catalogManager, schemaOperationDispatcher, entityStore,
idGenerator, secretManager);
- FunctionNormalizeDispatcher internalFunctionNormalizeDispatcher =
+ this.internalFunctionDispatcher =
new FunctionNormalizeDispatcher(functionOperationDispatcher,
catalogManager);
- this.internalFunctionDispatcher = internalFunctionNormalizeDispatcher;
- FunctionHookDispatcher functionHookDispatcher =
- new FunctionHookDispatcher(functionOperationDispatcher,
this::internalOwnerDispatcher);
- FunctionNormalizeDispatcher functionNormalizeDispatcher =
- new FunctionNormalizeDispatcher(functionHookDispatcher,
catalogManager);
- this.functionDispatcher = new FunctionEventDispatcher(eventBus,
functionNormalizeDispatcher);
+ return functionOperationDispatcher;
+ }
- // View operation chain: ViewEventDispatcher -> ViewNormalizeDispatcher ->
ViewHookDispatcher
- // -> ViewOperationDispatcher.
- ViewOperationDispatcher viewOperationDispatcher =
- new ViewOperationDispatcher(catalogManager, entityStore, idGenerator,
secretManager);
- this.internalViewDispatcher = viewOperationDispatcher;
+ private void initInternalViewDispatcher() {
ViewOperationDispatcher internalViewOperationDispatcher =
new ViewOperationDispatcher(
catalogManager,
@@ -974,12 +1001,9 @@ public class GravitinoEnv {
secretManager);
this.internalViewDispatcher =
new ViewNormalizeDispatcher(internalViewOperationDispatcher,
catalogManager);
- ViewHookDispatcher viewHookDispatcher =
- new ViewHookDispatcher(viewOperationDispatcher,
this::internalOwnerDispatcher);
- ViewNormalizeDispatcher viewNormalizeDispatcher =
- new ViewNormalizeDispatcher(viewHookDispatcher, catalogManager);
- this.viewDispatcher = new ViewEventDispatcher(eventBus,
viewNormalizeDispatcher);
+ }
+ private void initSemanticModelDispatcher(SchemaOperationDispatcher
schemaOperationDispatcher) {
// Semantic Model operation chain: SemanticModelNormalizeDispatcher ->
// SemanticModelOperationDispatcher -> ManagedSemanticModelOperations.
// TODO(#12595): Add Semantic Model event dispatching before server
integration.
@@ -989,48 +1013,43 @@ public class GravitinoEnv {
catalogManager, schemaOperationDispatcher, entityStore,
idGenerator, secretManager);
this.semanticModelDispatcher =
new SemanticModelNormalizeDispatcher(semanticModelOperationDispatcher,
catalogManager);
+ }
- this.statisticDispatcher =
- new StatisticEventDispatcher(
- eventBus, new StatisticManager(entityStore, idGenerator, config));
-
- // Create and initialize access control related modules
- boolean enableAuthorization = config.get(Configs.ENABLE_AUTHORIZATION);
- if (enableAuthorization) {
- AccessControlManager accessControlManager =
+ private void initInternalAuthorizationComponents() {
+ if (config.get(Configs.ENABLE_AUTHORIZATION)) {
+ this.internalAccessControlDispatcher =
new AccessControlManager(entityStore, idGenerator, config);
- this.internalAccessControlDispatcher = accessControlManager;
- AccessControlHookDispatcher accessControlHookDispatcher =
- new AccessControlHookDispatcher(accessControlManager);
- this.accessControlDispatcher =
- new AccessControlEventDispatcher(eventBus,
accessControlHookDispatcher);
- OwnerDispatcher ownerManager = new OwnerManager(entityStore);
- this.internalOwnerDispatcher = ownerManager;
- this.ownerDispatcher = new OwnerEventManager(eventBus, ownerManager);
+ this.internalOwnerDispatcher = new OwnerManager(entityStore);
this.bulkManager = new BulkManager(config);
- this.futureGrantManager = new FutureGrantManager(entityStore,
ownerManager);
+ this.futureGrantManager = new FutureGrantManager(entityStore,
internalOwnerDispatcher);
} else {
- this.accessControlDispatcher = null;
this.internalAccessControlDispatcher = null;
- this.ownerDispatcher = null;
this.internalOwnerDispatcher = null;
this.bulkManager = null;
this.futureGrantManager = null;
}
+ }
- this.auxServiceManager = new AuxiliaryServiceManager();
- this.auxServiceManager.serviceInit(config);
+ private void initEntityStoreAndCatalogManager() {
+ this.entityStore = EntityStoreFactory.createEntityStore(config);
+ entityStore.initialize(config);
- // Create and initialize Tag related modules
- TagManager tagManager = new TagManager(idGenerator, entityStore);
- this.internalTagDispatcher = tagManager;
- TagHookDispatcher tagHookDispatcher = new TagHookDispatcher(tagManager);
- this.tagDispatcher = new TagEventDispatcher(eventBus, tagHookDispatcher);
+ this.idGenerator = new RandomIdGenerator();
+ this.lockManager = new LockManager(config);
- PolicyManager policyManager = new PolicyManager(idGenerator, entityStore);
- this.internalPolicyDispatcher = policyManager;
- PolicyHookDispatcher policyHookDispatcher = new
PolicyHookDispatcher(policyManager);
- this.policyDispatcher = new PolicyEventDispatcher(eventBus,
policyHookDispatcher);
+ // CatalogManager must be initialized before MetalakeManager so force-drop
can remove child
+ // catalogs through CatalogManager.dropCatalog, the same path used by
FilesetCatalogOperations.
+ // CatalogManager registers its own change-log listener with compatible
entity stores, so no
+ // external poller wiring is needed here.
+ this.catalogManager = new CatalogManager(config, entityStore, idGenerator,
secretManager);
+ }
+
+ private void initGravitinoServerComponents() {
+ MetadataOperations metadataOperations = initMetadataComponents();
+ initPublicMetadataDispatchers(metadataOperations);
+
+ this.auxServiceManager = new AuxiliaryServiceManager();
+ this.auxServiceManager.serviceInit(config);
JobManager jobManager = new JobManager(config, entityStore, idGenerator);
JobTemplateValidationDispatcher validationDispatcher =
@@ -1045,4 +1064,119 @@ public class GravitinoEnv {
new BuiltInJobTemplateEventListener(jobManager, entityStore,
idGenerator);
eventListenerManager.addEventListener("builtin-job-template",
builtInJobTemplateListener);
}
+
+ private void initPublicMetadataDispatchers(MetadataOperations
metadataOperations) {
+ // Create and initialize metalake related modules, the operation chain is:
+ // MetalakeEventDispatcher -> MetalakeNormalizeDispatcher ->
MetalakeHookDispatcher ->
+ // MetalakeManager
+ MetalakeHookDispatcher metalakeHookDispatcher = new
MetalakeHookDispatcher(metalakeManager);
+ MetalakeNormalizeDispatcher metalakeNormalizeDispatcher =
+ new MetalakeNormalizeDispatcher(metalakeHookDispatcher);
+ this.metalakeDispatcher = new MetalakeEventDispatcher(eventBus,
metalakeNormalizeDispatcher);
+
+ // CatalogEventDispatcher -> CatalogNormalizeDispatcher ->
CatalogHookDispatcher ->
+ // CatalogManager
+ CatalogHookDispatcher catalogHookDispatcher = new
CatalogHookDispatcher(catalogManager);
+ CatalogNormalizeDispatcher catalogNormalizeDispatcher =
+ new CatalogNormalizeDispatcher(catalogHookDispatcher);
+ this.catalogDispatcher = new CatalogEventDispatcher(eventBus,
catalogNormalizeDispatcher);
+
+ FilesetHookDispatcher filesetHookDispatcher =
+ new
FilesetHookDispatcher(metadataOperations.filesetOperationDispatcher);
+ FilesetNormalizeDispatcher filesetNormalizeDispatcher =
+ new FilesetNormalizeDispatcher(filesetHookDispatcher, catalogManager);
+ this.filesetDispatcher = new FilesetEventDispatcher(eventBus,
filesetNormalizeDispatcher);
+
+ SchemaHookDispatcher schemaHookDispatcher =
+ new SchemaHookDispatcher(metadataOperations.schemaOperationDispatcher);
+ SchemaNormalizeDispatcher schemaNormalizeDispatcher =
+ new SchemaNormalizeDispatcher(schemaHookDispatcher, catalogManager);
+ this.schemaDispatcher = new SchemaEventDispatcher(eventBus,
schemaNormalizeDispatcher);
+
+ TableOperationDispatcher tableOperationDispatcher =
+ new TableOperationDispatcher(catalogManager, entityStore, idGenerator,
secretManager);
+ TableHookDispatcher tableHookDispatcher =
+ new TableHookDispatcher(tableOperationDispatcher,
this::internalOwnerDispatcher);
+ TableNormalizeDispatcher tableNormalizeDispatcher =
+ new TableNormalizeDispatcher(tableHookDispatcher, catalogManager);
+ this.tableDispatcher = new TableEventDispatcher(eventBus,
tableNormalizeDispatcher);
+
+ // TODO: We can install hooks when we need, we only supports ownership
post hook,
+ // partition doesn't have ownership, so we don't need it now.
+ this.partitionDispatcher = new PartitionEventDispatcher(eventBus,
internalPartitionDispatcher);
+
+ TopicHookDispatcher topicHookDispatcher =
+ new TopicHookDispatcher(metadataOperations.topicOperationDispatcher);
+ TopicNormalizeDispatcher topicNormalizeDispatcher =
+ new TopicNormalizeDispatcher(topicHookDispatcher, catalogManager);
+ this.topicDispatcher = new TopicEventDispatcher(eventBus,
topicNormalizeDispatcher);
+
+ ModelHookDispatcher modelHookDispatcher =
+ new ModelHookDispatcher(metadataOperations.modelOperationDispatcher);
+ ModelNormalizeDispatcher modelNormalizeDispatcher =
+ new ModelNormalizeDispatcher(modelHookDispatcher, catalogManager);
+ this.modelDispatcher = new ModelEventDispatcher(eventBus,
modelNormalizeDispatcher);
+
+ // Create and initialize Function related modules, the operation chain is:
+ // FunctionEventDispatcher -> FunctionNormalizeDispatcher ->
FunctionHookDispatcher ->
+ // FunctionOperationDispatcher
+ FunctionHookDispatcher functionHookDispatcher =
+ new FunctionHookDispatcher(
+ metadataOperations.functionOperationDispatcher,
this::internalOwnerDispatcher);
+ FunctionNormalizeDispatcher functionNormalizeDispatcher =
+ new FunctionNormalizeDispatcher(functionHookDispatcher,
catalogManager);
+ this.functionDispatcher = new FunctionEventDispatcher(eventBus,
functionNormalizeDispatcher);
+
+ // View operation chain: ViewEventDispatcher -> ViewNormalizeDispatcher ->
ViewHookDispatcher
+ // -> ViewOperationDispatcher.
+ ViewOperationDispatcher viewOperationDispatcher =
+ new ViewOperationDispatcher(catalogManager, entityStore, idGenerator,
secretManager);
+ ViewHookDispatcher viewHookDispatcher =
+ new ViewHookDispatcher(viewOperationDispatcher,
this::internalOwnerDispatcher);
+ ViewNormalizeDispatcher viewNormalizeDispatcher =
+ new ViewNormalizeDispatcher(viewHookDispatcher, catalogManager);
+ this.viewDispatcher = new ViewEventDispatcher(eventBus,
viewNormalizeDispatcher);
+
+ this.statisticDispatcher = new StatisticEventDispatcher(eventBus,
internalStatisticDispatcher);
+
+ // Create and initialize access control related modules
+ if (internalAccessControlDispatcher != null) {
+ AccessControlHookDispatcher accessControlHookDispatcher =
+ new AccessControlHookDispatcher(internalAccessControlDispatcher);
+ this.accessControlDispatcher =
+ new AccessControlEventDispatcher(eventBus,
accessControlHookDispatcher);
+ this.ownerDispatcher = new OwnerEventManager(eventBus,
internalOwnerDispatcher);
+ } else {
+ this.accessControlDispatcher = null;
+ this.ownerDispatcher = null;
+ }
+
+ // Create and initialize Tag related modules
+ TagHookDispatcher tagHookDispatcher = new
TagHookDispatcher(internalTagDispatcher);
+ this.tagDispatcher = new TagEventDispatcher(eventBus, tagHookDispatcher);
+
+ PolicyHookDispatcher policyHookDispatcher = new
PolicyHookDispatcher(internalPolicyDispatcher);
+ this.policyDispatcher = new PolicyEventDispatcher(eventBus,
policyHookDispatcher);
+ }
+
+ private static final class MetadataOperations {
+ private final FilesetOperationDispatcher filesetOperationDispatcher;
+ private final SchemaOperationDispatcher schemaOperationDispatcher;
+ private final TopicOperationDispatcher topicOperationDispatcher;
+ private final ModelOperationDispatcher modelOperationDispatcher;
+ private final FunctionOperationDispatcher functionOperationDispatcher;
+
+ private MetadataOperations(
+ FilesetOperationDispatcher filesetOperationDispatcher,
+ SchemaOperationDispatcher schemaOperationDispatcher,
+ TopicOperationDispatcher topicOperationDispatcher,
+ ModelOperationDispatcher modelOperationDispatcher,
+ FunctionOperationDispatcher functionOperationDispatcher) {
+ this.filesetOperationDispatcher = filesetOperationDispatcher;
+ this.schemaOperationDispatcher = schemaOperationDispatcher;
+ this.topicOperationDispatcher = topicOperationDispatcher;
+ this.modelOperationDispatcher = modelOperationDispatcher;
+ this.functionOperationDispatcher = functionOperationDispatcher;
+ }
+ }
}
diff --git
a/core/src/test/java/org/apache/gravitino/TestGravitinoEnvMetadataComponents.java
b/core/src/test/java/org/apache/gravitino/TestGravitinoEnvMetadataComponents.java
new file mode 100644
index 0000000000..ebffb1696b
--- /dev/null
+++
b/core/src/test/java/org/apache/gravitino/TestGravitinoEnvMetadataComponents.java
@@ -0,0 +1,336 @@
+/*
+ * 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.gravitino;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.Mockito.clearInvocations;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.withSettings;
+
+import java.lang.reflect.Field;
+import java.lang.reflect.Modifier;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.Map;
+import org.apache.commons.lang3.reflect.FieldUtils;
+import org.apache.gravitino.Entity.EntityType;
+import org.apache.gravitino.catalog.FilesetNormalizeDispatcher;
+import org.apache.gravitino.catalog.FilesetOperationDispatcher;
+import org.apache.gravitino.catalog.FunctionNormalizeDispatcher;
+import org.apache.gravitino.catalog.FunctionOperationDispatcher;
+import org.apache.gravitino.catalog.ModelNormalizeDispatcher;
+import org.apache.gravitino.catalog.ModelOperationDispatcher;
+import org.apache.gravitino.catalog.PartitionNormalizeDispatcher;
+import org.apache.gravitino.catalog.PartitionOperationDispatcher;
+import org.apache.gravitino.catalog.SchemaNormalizeDispatcher;
+import org.apache.gravitino.catalog.SchemaOperationDispatcher;
+import org.apache.gravitino.catalog.TableNormalizeDispatcher;
+import org.apache.gravitino.catalog.TableOperationDispatcher;
+import org.apache.gravitino.catalog.TopicNormalizeDispatcher;
+import org.apache.gravitino.catalog.TopicOperationDispatcher;
+import org.apache.gravitino.catalog.ViewNormalizeDispatcher;
+import org.apache.gravitino.catalog.ViewOperationDispatcher;
+import org.apache.gravitino.hook.FilesetHookDispatcher;
+import org.apache.gravitino.hook.FunctionHookDispatcher;
+import org.apache.gravitino.hook.ModelHookDispatcher;
+import org.apache.gravitino.hook.SchemaHookDispatcher;
+import org.apache.gravitino.hook.TableHookDispatcher;
+import org.apache.gravitino.hook.TopicHookDispatcher;
+import org.apache.gravitino.hook.ViewHookDispatcher;
+import org.apache.gravitino.listener.FilesetEventDispatcher;
+import org.apache.gravitino.listener.FunctionEventDispatcher;
+import org.apache.gravitino.listener.ModelEventDispatcher;
+import org.apache.gravitino.listener.PartitionEventDispatcher;
+import org.apache.gravitino.listener.SchemaEventDispatcher;
+import org.apache.gravitino.listener.StatisticEventDispatcher;
+import org.apache.gravitino.listener.TableEventDispatcher;
+import org.apache.gravitino.listener.TopicEventDispatcher;
+import org.apache.gravitino.listener.ViewEventDispatcher;
+import org.apache.gravitino.meta.BaseMetalake;
+import org.apache.gravitino.stats.StatisticManager;
+import org.apache.gravitino.stats.storage.MemoryPartitionStatsStorageFactory;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+
+class TestGravitinoEnvMetadataComponents {
+
+ private Map<Field, Object> singletonState;
+
+ @BeforeEach
+ void isolateSingletonState() throws IllegalAccessException {
+ GravitinoEnv singleton = GravitinoEnv.getInstance();
+ singletonState = snapshotState(singleton);
+ restoreState(singleton, snapshotState(new TestGravitinoEnv()));
+ }
+
+ @AfterEach
+ void restoreSingletonState() throws IllegalAccessException {
+ restoreState(GravitinoEnv.getInstance(), singletonState);
+ }
+
+ @Test
+ void testMetadataProfileRejectsNonSingletonEnvironment() {
+ IllegalStateException exception =
+ assertThrows(
+ IllegalStateException.class,
+ () -> new
TestGravitinoEnv().initializeMetadataComponents(metadataConfig(false)));
+
+ assertEquals(
+ "Metadata components must be initialized on
GravitinoEnv.getInstance().",
+ exception.getMessage());
+ }
+
+ @Test
+ void testInternalPartitionAndStatisticDispatchersRequireInitialization() {
+ GravitinoEnv env = GravitinoEnv.getInstance();
+
+ assertThrows(IllegalArgumentException.class,
env::internalPartitionDispatcher);
+ assertThrows(IllegalArgumentException.class,
env::internalStatisticDispatcher);
+ }
+
+ @Test
+ void
testMetadataProfileProvidesCompleteMetadataAccessWithoutServerServices() throws
Exception {
+ Config config = metadataConfig(false);
+ EntityStore entityStore = relationStore();
+ GravitinoEnv env = GravitinoEnv.getInstance();
+
+ try (MockedStatic<EntityStoreFactory> entityStoreFactory =
+ mockStatic(EntityStoreFactory.class)) {
+ entityStoreFactory
+ .when(() -> EntityStoreFactory.createEntityStore(config))
+ .thenReturn(entityStore);
+
+ env.initializeMetadataComponents(config);
+
+ assertSame(entityStore, env.entityStore());
+ assertNotNull(env.catalogManager());
+ assertNotNull(env.internalMetalakeDispatcher());
+ assertNotNull(env.internalCatalogDispatcher());
+ assertNotNull(env.internalFilesetDispatcher());
+ assertNotNull(env.internalSchemaDispatcher());
+ assertNotNull(env.internalTableDispatcher());
+ assertNotNull(env.internalPartitionDispatcher());
+ assertNotNull(env.internalTopicDispatcher());
+ assertNotNull(env.internalModelDispatcher());
+ assertNotNull(env.internalFunctionDispatcher());
+ assertNotNull(env.internalViewDispatcher());
+ assertNotNull(env.semanticModelDispatcher());
+ assertNotNull(env.credentialOperationDispatcher());
+ assertNotNull(env.secretPropertyOperationDispatcher());
+ assertNotNull(env.internalTagDispatcher());
+ assertNotNull(env.internalPolicyDispatcher());
+ assertInstanceOf(StatisticManager.class,
env.internalStatisticDispatcher());
+ assertNotNull(env.lockManager());
+ assertNotNull(env.metricsSystem());
+ assertNotNull(env.secretManager());
+
+ assertNull(env.auxServiceManager());
+ assertNull(env.eventListenerManager());
+ assertNull(env.metalakeDispatcher());
+ assertNull(env.catalogDispatcher());
+ assertNull(env.filesetDispatcher());
+ assertNull(env.schemaDispatcher());
+ assertNull(env.tableDispatcher());
+ assertNull(env.partitionDispatcher());
+ assertNull(env.topicDispatcher());
+ assertNull(env.modelDispatcher());
+ assertNull(env.functionDispatcher());
+ assertNull(env.viewDispatcher());
+ assertNull(env.tagDispatcher());
+ assertNull(env.policyDispatcher());
+ assertNull(env.statisticDispatcher());
+ assertNull(env.internalAccessControlDispatcher());
+ assertNull(env.internalOwnerDispatcher());
+ assertNull(env.bulkManager());
+ assertNull(env.futureGrantManager());
+ assertThrows(IllegalArgumentException.class, env::eventBus);
+ assertThrows(IllegalArgumentException.class,
env::jobOperationDispatcher);
+ assertThrows(IllegalArgumentException.class,
env::internalJobOperationDispatcher);
+
+ verify(entityStore).initialize(config);
+ clearInvocations(entityStore);
+ assertEquals(0, env.internalMetalakeDispatcher().listMetalakes().length);
+ verify(entityStore).list(Namespace.empty(), BaseMetalake.class,
EntityType.METALAKE);
+
+ assertDoesNotThrow(env::start);
+ } finally {
+ env.shutdown();
+ }
+
+ verify(entityStore).close();
+ }
+
+ @Test
+ void testMetadataProfileProvidesInternalAuthorizationWhenEnabled() throws
Exception {
+ Config config = metadataConfig(true);
+ EntityStore entityStore = relationStore();
+ GravitinoEnv env = GravitinoEnv.getInstance();
+
+ try (MockedStatic<EntityStoreFactory> entityStoreFactory =
+ mockStatic(EntityStoreFactory.class)) {
+ entityStoreFactory
+ .when(() -> EntityStoreFactory.createEntityStore(config))
+ .thenReturn(entityStore);
+
+ env.initializeMetadataComponents(config);
+
+ assertNotNull(env.internalAccessControlDispatcher());
+ assertNotNull(env.internalOwnerDispatcher());
+ assertNotNull(env.bulkManager());
+ assertNotNull(env.futureGrantManager());
+ assertNull(env.accessControlDispatcher());
+ assertNull(env.ownerDispatcher());
+ } finally {
+ env.shutdown();
+ }
+
+ verify(entityStore).close();
+ }
+
+ @Test
+ void testFullProfilePreservesDispatcherChains() throws Exception {
+ Config config = metadataConfig(false);
+ EntityStore entityStore = relationStore();
+ GravitinoEnv env = GravitinoEnv.getInstance();
+
+ try (MockedStatic<EntityStoreFactory> entityStoreFactory =
+ mockStatic(EntityStoreFactory.class)) {
+ entityStoreFactory
+ .when(() -> EntityStoreFactory.createEntityStore(config))
+ .thenReturn(entityStore);
+
+ env.initializeFullComponents(config);
+
+ assertDispatcherChain(
+ env.filesetDispatcher(),
+ FilesetEventDispatcher.class,
+ FilesetNormalizeDispatcher.class,
+ FilesetHookDispatcher.class,
+ FilesetOperationDispatcher.class);
+ assertDispatcherChain(
+ env.schemaDispatcher(),
+ SchemaEventDispatcher.class,
+ SchemaNormalizeDispatcher.class,
+ SchemaHookDispatcher.class,
+ SchemaOperationDispatcher.class);
+ assertDispatcherChain(
+ env.tableDispatcher(),
+ TableEventDispatcher.class,
+ TableNormalizeDispatcher.class,
+ TableHookDispatcher.class,
+ TableOperationDispatcher.class);
+ assertDispatcherChain(
+ env.topicDispatcher(),
+ TopicEventDispatcher.class,
+ TopicNormalizeDispatcher.class,
+ TopicHookDispatcher.class,
+ TopicOperationDispatcher.class);
+ assertDispatcherChain(
+ env.modelDispatcher(),
+ ModelEventDispatcher.class,
+ ModelNormalizeDispatcher.class,
+ ModelHookDispatcher.class,
+ ModelOperationDispatcher.class);
+ assertDispatcherChain(
+ env.functionDispatcher(),
+ FunctionEventDispatcher.class,
+ FunctionNormalizeDispatcher.class,
+ FunctionHookDispatcher.class,
+ FunctionOperationDispatcher.class);
+ assertDispatcherChain(
+ env.viewDispatcher(),
+ ViewEventDispatcher.class,
+ ViewNormalizeDispatcher.class,
+ ViewHookDispatcher.class,
+ ViewOperationDispatcher.class);
+ assertDispatcherChain(
+ env.partitionDispatcher(),
+ PartitionEventDispatcher.class,
+ PartitionNormalizeDispatcher.class,
+ PartitionOperationDispatcher.class);
+ assertSame(
+ env.internalPartitionDispatcher(),
+ FieldUtils.readField(env.partitionDispatcher(), "dispatcher", true));
+ assertDispatcherChain(
+ env.statisticDispatcher(), StatisticEventDispatcher.class,
StatisticManager.class);
+ assertSame(
+ env.internalStatisticDispatcher(),
+ FieldUtils.readField(env.statisticDispatcher(), "dispatcher", true));
+ } finally {
+ env.shutdown();
+ }
+
+ verify(entityStore).close();
+ }
+
+ private static Config metadataConfig(boolean enableAuthorization) {
+ Config config = new Config(false) {};
+ config.set(Configs.ENABLE_AUTHORIZATION, enableAuthorization);
+ config.set(Configs.SERVICE_ADMINS, Collections.singletonList("admin"));
+ config.set(
+ Configs.PARTITION_STATS_STORAGE_FACTORY_CLASS,
+ MemoryPartitionStatsStorageFactory.class.getCanonicalName());
+ return config;
+ }
+
+ private static EntityStore relationStore() {
+ return mock(
+ EntityStore.class,
withSettings().extraInterfaces(SupportsRelationOperations.class));
+ }
+
+ private static Map<Field, Object> snapshotState(GravitinoEnv env) throws
IllegalAccessException {
+ Map<Field, Object> state = new LinkedHashMap<>();
+ for (Field field : FieldUtils.getAllFieldsList(GravitinoEnv.class)) {
+ if (!Modifier.isStatic(field.getModifiers())) {
+ state.put(field, FieldUtils.readField(field, env, true));
+ }
+ }
+ return state;
+ }
+
+ private static void restoreState(GravitinoEnv env, Map<Field, Object> state)
+ throws IllegalAccessException {
+ for (Map.Entry<Field, Object> entry : state.entrySet()) {
+ FieldUtils.writeField(entry.getKey(), env, entry.getValue(), true);
+ }
+ }
+
+ private static void assertDispatcherChain(Object dispatcher, Class<?>...
dispatcherClasses)
+ throws IllegalAccessException {
+ Object currentDispatcher = dispatcher;
+ for (int index = 0; index < dispatcherClasses.length; index++) {
+ assertInstanceOf(dispatcherClasses[index], currentDispatcher);
+ if (index < dispatcherClasses.length - 1) {
+ currentDispatcher = FieldUtils.readField(currentDispatcher,
"dispatcher", true);
+ }
+ }
+ }
+
+ private static final class TestGravitinoEnv extends GravitinoEnv {}
+}