This is an automated email from the ASF dual-hosted git repository.
roryqi 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 cfa855de40 [#12233] feat(storage): Store tag assignment values (#12354)
cfa855de40 is described below
commit cfa855de4098797bc2fe6179abc6c9f0e17396b5
Author: roryqi <[email protected]>
AuthorDate: Wed Aug 5 21:29:54 2026 +0800
[#12233] feat(storage): Store tag assignment values (#12354)
### What changes were proposed in this pull request?
This PR stores tag assignment values in the relational storage layer.
Changes include:
- Reuse `TagEntity` to expose assignment values in tag-object relation
context.
- Persist tag value constraints in tag metadata storage.
- Persist assignment values in `tag_metadata_object_rel`.
- Support listing associated metadata objects by exact tag assignment
value.
- Update H2/MySQL/PostgreSQL schema files and relational SQL providers.
### Why are the changes needed?
This is needed to support assigning multiple values to a tag and
retrieving assigned tag values from storage.
Fix: #12233
### Does this PR introduce _any_ user-facing change?
No new public API is introduced in this PR. It adds relational storage
support for tag value constraints and tag assignment values used by the
tag value APIs.
### How was this patch tested?
- `./gradlew :core:spotlessApply`
- `./gradlew :core:test --tests org.apache.gravitino.meta.TestTagEntity
--tests
org.apache.gravitino.storage.relational.service.TestTagMetaService
-PskipITs -PskipDockerTests=false`
---
.../java/org/apache/gravitino/meta/TagEntity.java | 70 ++++-
.../gravitino/storage/relational/JDBCBackend.java | 88 ++++++
.../mapper/TagMetadataObjectRelMapper.java | 18 +-
.../TagMetadataObjectRelSQLProviderFactory.java | 17 ++
.../provider/base/TagMetaBaseSQLProvider.java | 11 +-
.../base/TagMetadataObjectRelBaseSQLProvider.java | 50 +++-
.../postgresql/TagMetaPostgreSQLProvider.java | 5 +-
.../TagMetadataObjectRelPostgreSQLProvider.java | 21 +-
.../relational/po/TagMetadataObjectRelPO.java | 8 +
.../gravitino/storage/relational/po/TagPO.java | 16 ++
.../storage/relational/service/TagMetaService.java | 306 ++++++++++++++++++---
.../storage/relational/utils/POConverters.java | 35 +++
.../java/org/apache/gravitino/tag/TagManager.java | 2 +-
.../org/apache/gravitino/meta/TestTagEntity.java | 63 +++++
.../relational/service/TestTagMetaService.java | 234 ++++++++++++++--
.../storage/relational/utils/TestPOConverters.java | 1 +
scripts/h2/schema-2.0.0-h2.sql | 7 +-
scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql | 9 +
scripts/mysql/schema-2.0.0-mysql.sql | 7 +-
scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql | 12 +
scripts/postgresql/schema-2.0.0-postgresql.sql | 7 +-
.../upgrade-1.3.0-to-2.0.0-postgresql.sql | 11 +
22 files changed, 919 insertions(+), 79 deletions(-)
diff --git a/core/src/main/java/org/apache/gravitino/meta/TagEntity.java
b/core/src/main/java/org/apache/gravitino/meta/TagEntity.java
index 1eabad889a..2730f921c9 100644
--- a/core/src/main/java/org/apache/gravitino/meta/TagEntity.java
+++ b/core/src/main/java/org/apache/gravitino/meta/TagEntity.java
@@ -20,10 +20,12 @@
package org.apache.gravitino.meta;
import com.google.common.collect.Maps;
+import java.util.Arrays;
import java.util.Collections;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
+import javax.annotation.Nullable;
import org.apache.gravitino.Audit;
import org.apache.gravitino.Auditable;
import org.apache.gravitino.Entity;
@@ -31,6 +33,8 @@ import org.apache.gravitino.Field;
import org.apache.gravitino.HasIdentifier;
import org.apache.gravitino.Namespace;
import org.apache.gravitino.tag.Tag;
+import org.apache.gravitino.tag.TagAssignment;
+import org.apache.gravitino.tag.TagValueConstraint;
public class TagEntity implements Tag, Entity, Auditable, HasIdentifier {
@@ -46,6 +50,9 @@ public class TagEntity implements Tag, Entity, Auditable,
HasIdentifier {
public static final Field PROPERTIES =
Field.optional("properties", Map.class, "The properties of the tag
entity.");
+ public static final Field ALLOWED_VALUES =
+ Field.optional("allowed_values", String[].class, "The allowed assignment
values.");
+
public static final Field AUDIT_INFO =
Field.required("audit_info", Audit.class, "The audit details of the tag
entity.");
@@ -54,6 +61,8 @@ public class TagEntity implements Tag, Entity, Auditable,
HasIdentifier {
private Namespace namespace;
private String comment;
private Map<String, String> properties;
+ private String[] allowedValues;
+ @Nullable private TagAssignment assignment;
private Audit auditInfo;
private TagEntity() {}
@@ -65,6 +74,7 @@ public class TagEntity implements Tag, Entity, Auditable,
HasIdentifier {
fields.put(NAME, name);
fields.put(COMMENT, comment);
fields.put(PROPERTIES, properties);
+ fields.put(ALLOWED_VALUES, allowedValues);
fields.put(AUDIT_INFO, auditInfo);
return Collections.unmodifiableMap(fields);
@@ -100,6 +110,22 @@ public class TagEntity implements Tag, Entity, Auditable,
HasIdentifier {
return properties;
}
+ @Override
+ public TagValueConstraint valueConstraint() {
+ if (allowedValues == null) {
+ return TagValueConstraint.anyValue();
+ }
+ if (allowedValues.length == 0) {
+ return TagValueConstraint.noValue();
+ }
+ return TagValueConstraint.ofAllowedValues(allowedValues);
+ }
+
+ @Override
+ public Optional<TagAssignment> assignment() {
+ return Optional.ofNullable(assignment);
+ }
+
@Override
public Optional<Boolean> inherited() {
return Optional.empty();
@@ -125,12 +151,37 @@ public class TagEntity implements Tag, Entity, Auditable,
HasIdentifier {
&& Objects.equals(namespace, that.namespace)
&& Objects.equals(comment, that.comment)
&& Objects.equals(properties, that.properties)
+ && Arrays.equals(allowedValues, that.allowedValues)
&& Objects.equals(auditInfo, that.auditInfo);
}
@Override
public int hashCode() {
- return Objects.hash(id, name, namespace, comment, properties, auditInfo);
+ int result = Objects.hash(id, name, namespace, comment, properties,
auditInfo);
+ result = 31 * result + Arrays.hashCode(allowedValues);
+ return result;
+ }
+
+ /**
+ * Returns a copy of this tag entity with the given assignment context.
+ *
+ * <p>The assignment context is used only when the tag is returned for a
metadata object relation;
+ * it is not part of the tag definition fields.
+ *
+ * @param assignment The assignment context, or null to clear it.
+ * @return The copied tag entity with the assignment context.
+ */
+ public TagEntity copyWithAssignment(@Nullable TagAssignment assignment) {
+ return TagEntity.builder()
+ .withId(id)
+ .withName(name)
+ .withNamespace(namespace)
+ .withComment(comment)
+ .withProperties(properties)
+ .withAllowedValues(allowedValues)
+ .withAssignment(assignment)
+ .withAuditInfo(auditInfo)
+ .build();
}
public static Builder builder() {
@@ -170,6 +221,23 @@ public class TagEntity implements Tag, Entity, Auditable,
HasIdentifier {
return this;
}
+ public Builder withAllowedValues(String[] allowedValues) {
+ tagEntity.allowedValues = allowedValues == null ? null :
allowedValues.clone();
+ return this;
+ }
+
+ /**
+ * Sets the assignment context for this tag entity. This is only used when
the tag is returned
+ * for a metadata object relation.
+ *
+ * @param assignment The assignment context, or null to clear it.
+ * @return The builder instance.
+ */
+ public Builder withAssignment(@Nullable TagAssignment assignment) {
+ tagEntity.assignment = assignment;
+ return this;
+ }
+
public Builder withAuditInfo(Audit auditInfo) {
tagEntity.auditInfo = auditInfo;
return this;
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/JDBCBackend.java
b/core/src/main/java/org/apache/gravitino/storage/relational/JDBCBackend.java
index ada6e6160b..c87edc0c38 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/JDBCBackend.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/JDBCBackend.java
@@ -26,6 +26,7 @@ import com.google.common.base.Preconditions;
import com.google.common.collect.ImmutableMap;
import com.google.common.collect.Lists;
import java.io.IOException;
+import java.util.Arrays;
import java.util.List;
import java.util.Map;
import java.util.function.Function;
@@ -39,6 +40,9 @@ import org.apache.gravitino.HasIdentifier;
import org.apache.gravitino.MetadataObject;
import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.Namespace;
+import org.apache.gravitino.RelationEdgeTarget;
+import org.apache.gravitino.RelationQuery;
+import org.apache.gravitino.RelationUpdate;
import org.apache.gravitino.RelationalEntity;
import org.apache.gravitino.SupportsRelationOperations;
import org.apache.gravitino.UnsupportedEntityTypeException;
@@ -87,6 +91,7 @@ import
org.apache.gravitino.storage.relational.service.TopicMetaService;
import org.apache.gravitino.storage.relational.service.UserMetaService;
import org.apache.gravitino.storage.relational.service.ViewMetaService;
import org.apache.gravitino.storage.relational.session.SqlSessionFactoryHelper;
+import org.apache.gravitino.tag.TagValue;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -860,6 +865,89 @@ public class JDBCBackend implements RelationalBackend,
SupportsOrphanedRelationC
}
}
+ @Override
+ public <E extends Entity & HasIdentifier> List<E>
listEntitiesByRelation(RelationQuery query)
+ throws IOException {
+ if (!query.relationValue().isPresent()) {
+ return listEntitiesByRelation(
+ query.relationType(),
+ query.anchorIdentifier(),
+ query.anchorEntityType(),
+ query.allFields());
+ }
+
+ switch (query.relationType()) {
+ case TAG_METADATA_OBJECT_REL:
+ Preconditions.checkArgument(
+ query.anchorEntityType() == Entity.EntityType.TAG,
+ "Relation value filter is only supported when listing metadata
objects for a tag");
+ return (List<E>)
+ TagMetaService.getInstance()
+ .listAssociatedMetadataObjectsForTag(
+ query.anchorIdentifier(), query.relationValue().get());
+ default:
+ throw new IllegalArgumentException(
+ String.format(
+ "Relation value filter is not supported for relation type %s",
+ query.relationType()));
+ }
+ }
+
+ @Override
+ public <E extends Entity & HasIdentifier> List<E>
updateEntityRelations(RelationUpdate update)
+ throws IOException, NoSuchEntityException, EntityAlreadyExistsException {
+ switch (update.relationType()) {
+ case TAG_METADATA_OBJECT_REL:
+ return (List<E>)
+ TagMetaService.getInstance()
+ .associateTagValuesWithMetadataObject(
+ update.sourceIdentifier(),
+ update.sourceEntityType(),
+ toTagValues(update.targetsToAdd()),
+ toTagValues(update.targetsToRemove()));
+ default:
+ Preconditions.checkArgument(
+ !update.hasRelationValues(),
+ "Relation values are not supported for relation type %s",
+ update.relationType());
+ return updateEntityRelations(
+ update.relationType(),
+ update.sourceIdentifier(),
+ update.sourceEntityType(),
+ toNameIdentifiers(update.targetsToAdd()),
+ toNameIdentifiers(update.targetsToRemove()));
+ }
+ }
+
+ private static NameIdentifier[] toNameIdentifiers(RelationEdgeTarget[]
relationTargets) {
+ if (relationTargets == null) {
+ return null;
+ }
+
+ return Arrays.stream(relationTargets)
+ .map(RelationEdgeTarget::nameIdentifier)
+ .toArray(NameIdentifier[]::new);
+ }
+
+ private static TagValue[] toTagValues(RelationEdgeTarget[] relationTargets) {
+ if (relationTargets == null) {
+ return null;
+ }
+
+ return
Arrays.stream(relationTargets).map(JDBCBackend::toTagValue).toArray(TagValue[]::new);
+ }
+
+ private static TagValue toTagValue(RelationEdgeTarget relationTarget) {
+ Preconditions.checkArgument(
+ relationTarget.entityType() == Entity.EntityType.TAG,
+ "Relation target type must be TAG for tag metadata object relations,
but is %s",
+ relationTarget.entityType());
+ return relationTarget
+ .relationValue()
+ .map(value -> TagValue.of(relationTarget.nameIdentifier().name(),
value))
+ .orElseGet(() ->
TagValue.noValue(relationTarget.nameIdentifier().name()));
+ }
+
@Override
public <E extends Entity & HasIdentifier> E getEntityByRelation(
Type relType,
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelMapper.java
index 14d6d87f7b..bdb461a3f6 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelMapper.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelMapper.java
@@ -40,7 +40,7 @@ public interface TagMetadataObjectRelMapper {
@SelectProvider(
type = TagMetadataObjectRelSQLProviderFactory.class,
method = "getTagPOsByMetadataObjectAndTagName")
- TagPO getTagPOsByMetadataObjectAndTagName(
+ List<TagPO> getTagPOsByMetadataObjectAndTagName(
@Param("metadataObjectId") Long metadataObjectId,
@Param("metadataObjectType") String metadataObjectType,
@Param("tagName") String tagName);
@@ -51,6 +51,14 @@ public interface TagMetadataObjectRelMapper {
List<TagMetadataObjectRelPO> listTagMetadataObjectRelsByMetalakeAndTagName(
@Param("metalakeName") String metalakeName, @Param("tagName") String
tagName);
+ @SelectProvider(
+ type = TagMetadataObjectRelSQLProviderFactory.class,
+ method = "listTagMetadataObjectRelsByMetalakeAndTagNameAndValue")
+ List<TagMetadataObjectRelPO>
listTagMetadataObjectRelsByMetalakeAndTagNameAndValue(
+ @Param("metalakeName") String metalakeName,
+ @Param("tagName") String tagName,
+ @Param("tagValue") String tagValue);
+
@InsertProvider(
type = TagMetadataObjectRelSQLProviderFactory.class,
method = "batchInsertTagMetadataObjectRels")
@@ -64,6 +72,14 @@ public interface TagMetadataObjectRelMapper {
@Param("metadataObjectType") String metadataObjectType,
@Param("tagIds") List<Long> tagIds);
+ @UpdateProvider(
+ type = TagMetadataObjectRelSQLProviderFactory.class,
+ method =
"batchDeleteTagMetadataObjectRelsByTagIdsAndValuesAndMetadataObject")
+ void batchDeleteTagMetadataObjectRelsByTagIdsAndValuesAndMetadataObject(
+ @Param("metadataObjectId") Long metadataObjectId,
+ @Param("metadataObjectType") String metadataObjectType,
+ @Param("tagRels") List<TagMetadataObjectRelPO> tagRelPOs);
+
@UpdateProvider(
type = TagMetadataObjectRelSQLProviderFactory.class,
method = "softDeleteTagMetadataObjectRelsByMetalakeAndTagName")
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelSQLProviderFactory.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelSQLProviderFactory.java
index fd2f49e645..43ec136a93 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelSQLProviderFactory.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelSQLProviderFactory.java
@@ -71,6 +71,14 @@ public class TagMetadataObjectRelSQLProviderFactory {
return
getProvider().listTagMetadataObjectRelsByMetalakeAndTagName(metalakeName,
tagName);
}
+ public static String listTagMetadataObjectRelsByMetalakeAndTagNameAndValue(
+ @Param("metalakeName") String metalakeName,
+ @Param("tagName") String tagName,
+ @Param("tagValue") String tagValue) {
+ return getProvider()
+ .listTagMetadataObjectRelsByMetalakeAndTagNameAndValue(metalakeName,
tagName, tagValue);
+ }
+
public static String batchInsertTagMetadataObjectRels(
@Param("tagRels") List<TagMetadataObjectRelPO> tagRelPOs) {
return getProvider().batchInsertTagMetadataObjectRels(tagRelPOs);
@@ -85,6 +93,15 @@ public class TagMetadataObjectRelSQLProviderFactory {
metadataObjectId, metadataObjectType, tagIds);
}
+ public static String
batchDeleteTagMetadataObjectRelsByTagIdsAndValuesAndMetadataObject(
+ @Param("metadataObjectId") Long metadataObjectId,
+ @Param("metadataObjectType") String metadataObjectType,
+ @Param("tagRels") List<TagMetadataObjectRelPO> tagRelPOs) {
+ return getProvider()
+ .batchDeleteTagMetadataObjectRelsByTagIdsAndValuesAndMetadataObject(
+ metadataObjectId, metadataObjectType, tagRelPOs);
+ }
+
public static String softDeleteTagMetadataObjectRelsByMetalakeAndTagName(
@Param("metalakeName") String metalakeName, @Param("tagName") String
tagName) {
return
getProvider().softDeleteTagMetadataObjectRelsByMetalakeAndTagName(metalakeName,
tagName);
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TagMetaBaseSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TagMetaBaseSQLProvider.java
index 8aac66bf88..e33f2080c9 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TagMetaBaseSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TagMetaBaseSQLProvider.java
@@ -32,6 +32,7 @@ public class TagMetaBaseSQLProvider {
+ " tm.metalake_id as metalakeId,"
+ " tm.tag_comment as comment,"
+ " tm.properties as properties,"
+ + " tm.allowed_values as allowedValues,"
+ " tm.audit_info as auditInfo,"
+ " tm.current_version as currentVersion,"
+ " tm.last_version as lastVersion,"
@@ -51,6 +52,7 @@ public class TagMetaBaseSQLProvider {
+ " tm.metalake_id as metalakeId,"
+ " tm.tag_comment as comment,"
+ " tm.properties as properties,"
+ + " tm.allowed_values as allowedValues,"
+ " tm.audit_info as auditInfo,"
+ " tm.current_version as currentVersion,"
+ " tm.last_version as lastVersion,"
@@ -86,6 +88,7 @@ public class TagMetaBaseSQLProvider {
+ " tm.metalake_id as metalakeId,"
+ " tm.tag_comment as comment,"
+ " tm.properties as properties,"
+ + " tm.allowed_values as allowedValues,"
+ " tm.audit_info as auditInfo,"
+ " tm.current_version as currentVersion,"
+ " tm.last_version as lastVersion,"
@@ -103,7 +106,7 @@ public class TagMetaBaseSQLProvider {
return "INSERT INTO "
+ TAG_TABLE_NAME
+ " (tag_id, tag_name,"
- + " metalake_id, tag_comment, properties, audit_info,"
+ + " metalake_id, tag_comment, properties, allowed_values, audit_info,"
+ " current_version, last_version, deleted_at)"
+ " VALUES ("
+ " #{tagMeta.tagId},"
@@ -111,6 +114,7 @@ public class TagMetaBaseSQLProvider {
+ " #{tagMeta.metalakeId},"
+ " #{tagMeta.comment},"
+ " #{tagMeta.properties},"
+ + " #{tagMeta.allowedValues},"
+ " #{tagMeta.auditInfo},"
+ " #{tagMeta.currentVersion},"
+ " #{tagMeta.lastVersion},"
@@ -122,7 +126,7 @@ public class TagMetaBaseSQLProvider {
return "INSERT INTO "
+ TAG_TABLE_NAME
+ " (tag_id, tag_name,"
- + " metalake_id, tag_comment, properties, audit_info,"
+ + " metalake_id, tag_comment, properties, allowed_values, audit_info,"
+ " current_version, last_version, deleted_at)"
+ " VALUES ("
+ " #{tagMeta.tagId},"
@@ -130,6 +134,7 @@ public class TagMetaBaseSQLProvider {
+ " #{tagMeta.metalakeId},"
+ " #{tagMeta.comment},"
+ " #{tagMeta.properties},"
+ + " #{tagMeta.allowedValues},"
+ " #{tagMeta.auditInfo},"
+ " #{tagMeta.currentVersion},"
+ " #{tagMeta.lastVersion},"
@@ -140,6 +145,7 @@ public class TagMetaBaseSQLProvider {
+ " metalake_id = #{tagMeta.metalakeId},"
+ " tag_comment = #{tagMeta.comment},"
+ " properties = #{tagMeta.properties},"
+ + " allowed_values = #{tagMeta.allowedValues},"
+ " audit_info = #{tagMeta.auditInfo},"
+ " current_version = #{tagMeta.currentVersion},"
+ " last_version = #{tagMeta.lastVersion},"
@@ -153,6 +159,7 @@ public class TagMetaBaseSQLProvider {
+ " SET tag_name = #{newTagMeta.tagName},"
+ " tag_comment = #{newTagMeta.comment},"
+ " properties = #{newTagMeta.properties},"
+ + " allowed_values = #{newTagMeta.allowedValues},"
+ " audit_info = #{newTagMeta.auditInfo},"
+ " current_version = #{newTagMeta.currentVersion},"
+ " last_version = #{newTagMeta.lastVersion},"
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TagMetadataObjectRelBaseSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TagMetadataObjectRelBaseSQLProvider.java
index 1f3a727066..c1185b96d5 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TagMetadataObjectRelBaseSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TagMetadataObjectRelBaseSQLProvider.java
@@ -41,6 +41,8 @@ public class TagMetadataObjectRelBaseSQLProvider {
@Param("metadataObjectType") String metadataObjectType) {
return "SELECT tm.tag_id as tagId, tm.tag_name as tagName,"
+ " tm.metalake_id as metalakeId, tm.tag_comment as comment,
tm.properties as properties,"
+ + " tm.allowed_values as allowedValues,"
+ + " NULLIF(te.tag_value, '') as assignmentValue,"
+ " tm.audit_info as auditInfo,"
+ " tm.current_version as currentVersion,"
+ " tm.last_version as lastVersion,"
@@ -61,6 +63,8 @@ public class TagMetadataObjectRelBaseSQLProvider {
@Param("tagName") String tagName) {
return "SELECT tm.tag_id as tagId, tm.tag_name as tagName,"
+ " tm.metalake_id as metalakeId, tm.tag_comment as comment,
tm.properties as properties,"
+ + " tm.allowed_values as allowedValues,"
+ + " NULLIF(te.tag_value, '') as assignmentValue,"
+ " tm.audit_info as auditInfo,"
+ " tm.current_version as currentVersion,"
+ " tm.last_version as lastVersion,"
@@ -78,7 +82,8 @@ public class TagMetadataObjectRelBaseSQLProvider {
public String listTagMetadataObjectRelsByMetalakeAndTagName(
@Param("metalakeName") String metalakeName, @Param("tagName") String
tagName) {
return "SELECT te.tag_id as tagId, te.metadata_object_id as
metadataObjectId,"
- + " te.metadata_object_type as metadataObjectType, te.audit_info as
auditInfo,"
+ + " te.metadata_object_type as metadataObjectType, te.tag_value as
tagValue,"
+ + " te.audit_info as auditInfo,"
+ " te.current_version as currentVersion, te.last_version as
lastVersion,"
+ " te.deleted_at as deletedAt"
+ " FROM "
@@ -92,18 +97,40 @@ public class TagMetadataObjectRelBaseSQLProvider {
+ " AND te.deleted_at = 0 AND tm.deleted_at = 0 AND mm.deleted_at = 0";
}
+ public String listTagMetadataObjectRelsByMetalakeAndTagNameAndValue(
+ @Param("metalakeName") String metalakeName,
+ @Param("tagName") String tagName,
+ @Param("tagValue") String tagValue) {
+ return "SELECT te.tag_id as tagId, te.metadata_object_id as
metadataObjectId,"
+ + " te.metadata_object_type as metadataObjectType, te.tag_value as
tagValue,"
+ + " te.audit_info as auditInfo,"
+ + " te.current_version as currentVersion, te.last_version as
lastVersion,"
+ + " te.deleted_at as deletedAt"
+ + " FROM "
+ + TagMetadataObjectRelMapper.TAG_METADATA_OBJECT_RELATION_TABLE_NAME
+ + " te JOIN "
+ + TagMetaMapper.TAG_TABLE_NAME
+ + " tm ON te.tag_id = tm.tag_id JOIN "
+ + MetalakeMetaMapper.TABLE_NAME
+ + " mm ON tm.metalake_id = mm.metalake_id"
+ + " WHERE mm.metalake_name = #{metalakeName} AND tm.tag_name =
#{tagName}"
+ + " AND te.tag_value = #{tagValue}"
+ + " AND te.deleted_at = 0 AND tm.deleted_at = 0 AND mm.deleted_at = 0";
+ }
+
public String batchInsertTagMetadataObjectRels(
@Param("tagRels") List<TagMetadataObjectRelPO> tagRelPOs) {
return "<script>"
+ "INSERT INTO "
+ TagMetadataObjectRelMapper.TAG_METADATA_OBJECT_RELATION_TABLE_NAME
- + " (tag_id, metadata_object_id, metadata_object_type, audit_info,"
+ + " (tag_id, metadata_object_id, metadata_object_type, tag_value,
audit_info,"
+ " current_version, last_version, deleted_at)"
+ " VALUES "
+ "<foreach collection='tagRels' item='item' separator=','>"
+ "(#{item.tagId},"
+ " #{item.metadataObjectId},"
+ " #{item.metadataObjectType},"
+ + " #{item.tagValue},"
+ " #{item.auditInfo},"
+ " #{item.currentVersion},"
+ " #{item.lastVersion},"
@@ -130,6 +157,25 @@ public class TagMetadataObjectRelBaseSQLProvider {
+ "</script>";
}
+ public String
batchDeleteTagMetadataObjectRelsByTagIdsAndValuesAndMetadataObject(
+ @Param("metadataObjectId") Long metadataObjectId,
+ @Param("metadataObjectType") String metadataObjectType,
+ @Param("tagRels") List<TagMetadataObjectRelPO> tagRelPOs) {
+ return "<script>"
+ + "UPDATE "
+ + TagMetadataObjectRelMapper.TAG_METADATA_OBJECT_RELATION_TABLE_NAME
+ + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)"
+ + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000"
+ + " WHERE metadata_object_id = #{metadataObjectId}"
+ + " AND metadata_object_type = #{metadataObjectType} AND deleted_at =
0"
+ + " AND ("
+ + "<foreach item='item' collection='tagRels' separator=' OR '>"
+ + "(tag_id = #{item.tagId} AND tag_value = #{item.tagValue})"
+ + "</foreach>"
+ + ")"
+ + "</script>";
+ }
+
public String softDeleteTagMetadataObjectRelsByMetalakeAndTagName(
@Param("metalakeName") String metalakeName, @Param("tagName") String
tagName) {
return "UPDATE "
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TagMetaPostgreSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TagMetaPostgreSQLProvider.java
index 785e124d2c..cf828a67c5 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TagMetaPostgreSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TagMetaPostgreSQLProvider.java
@@ -51,7 +51,7 @@ public class TagMetaPostgreSQLProvider extends
TagMetaBaseSQLProvider {
return "INSERT INTO "
+ TAG_TABLE_NAME
+ " (tag_id, tag_name,"
- + " metalake_id, tag_comment, properties, audit_info,"
+ + " metalake_id, tag_comment, properties, allowed_values, audit_info,"
+ " current_version, last_version, deleted_at)"
+ " VALUES ("
+ " #{tagMeta.tagId},"
@@ -59,6 +59,7 @@ public class TagMetaPostgreSQLProvider extends
TagMetaBaseSQLProvider {
+ " #{tagMeta.metalakeId},"
+ " #{tagMeta.comment},"
+ " #{tagMeta.properties},"
+ + " #{tagMeta.allowedValues},"
+ " #{tagMeta.auditInfo},"
+ " #{tagMeta.currentVersion},"
+ " #{tagMeta.lastVersion},"
@@ -69,6 +70,7 @@ public class TagMetaPostgreSQLProvider extends
TagMetaBaseSQLProvider {
+ " metalake_id = #{tagMeta.metalakeId},"
+ " tag_comment = #{tagMeta.comment},"
+ " properties = #{tagMeta.properties},"
+ + " allowed_values = #{tagMeta.allowedValues},"
+ " audit_info = #{tagMeta.auditInfo},"
+ " current_version = #{tagMeta.currentVersion},"
+ " last_version = #{tagMeta.lastVersion},"
@@ -83,6 +85,7 @@ public class TagMetaPostgreSQLProvider extends
TagMetaBaseSQLProvider {
+ " SET tag_name = #{newTagMeta.tagName},"
+ " tag_comment = #{newTagMeta.comment},"
+ " properties = #{newTagMeta.properties},"
+ + " allowed_values = #{newTagMeta.allowedValues},"
+ " audit_info = #{newTagMeta.auditInfo},"
+ " current_version = #{newTagMeta.currentVersion},"
+ " last_version = #{newTagMeta.lastVersion},"
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TagMetadataObjectRelPostgreSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TagMetadataObjectRelPostgreSQLProvider.java
index 992e105ee6..b8a2231aa5 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TagMetadataObjectRelPostgreSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TagMetadataObjectRelPostgreSQLProvider.java
@@ -34,9 +34,27 @@ import
org.apache.gravitino.storage.relational.mapper.TagMetadataObjectRelMapper
import org.apache.gravitino.storage.relational.mapper.TopicMetaMapper;
import org.apache.gravitino.storage.relational.mapper.ViewMetaMapper;
import
org.apache.gravitino.storage.relational.mapper.provider.base.TagMetadataObjectRelBaseSQLProvider;
+import org.apache.gravitino.storage.relational.po.TagMetadataObjectRelPO;
import org.apache.ibatis.annotations.Param;
public class TagMetadataObjectRelPostgreSQLProvider extends
TagMetadataObjectRelBaseSQLProvider {
+ @Override
+ public String
batchDeleteTagMetadataObjectRelsByTagIdsAndValuesAndMetadataObject(
+ Long metadataObjectId, String metadataObjectType,
List<TagMetadataObjectRelPO> tagRelPOs) {
+ return "<script>"
+ + "UPDATE "
+ + TAG_METADATA_OBJECT_RELATION_TABLE_NAME
+ + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000
AS BIGINT)"
+ + " WHERE metadata_object_id = #{metadataObjectId}"
+ + " AND metadata_object_type = #{metadataObjectType} AND deleted_at =
0"
+ + " AND ("
+ + "<foreach item='item' collection='tagRels' separator=' OR '>"
+ + "(tag_id = #{item.tagId} AND tag_value = #{item.tagValue})"
+ + "</foreach>"
+ + ")"
+ + "</script>";
+ }
+
@Override
public String softDeleteTagMetadataObjectRelsByMetalakeAndTagName(
String metalakeName, String tagName) {
@@ -238,7 +256,8 @@ public class TagMetadataObjectRelPostgreSQLProvider extends
TagMetadataObjectRel
@Override
public String listTagMetadataObjectRelsByMetalakeAndTagName(String
metalakeName, String tagName) {
return "SELECT te.tag_id as tagId, te.metadata_object_id as
metadataObjectId,"
- + " te.metadata_object_type as metadataObjectType, te.audit_info as
auditInfo,"
+ + " te.metadata_object_type as metadataObjectType, te.tag_value as
tagValue,"
+ + " te.audit_info as auditInfo,"
+ " te.current_version as currentVersion, te.last_version as
lastVersion,"
+ " te.deleted_at as deletedAt"
+ " FROM "
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/po/TagMetadataObjectRelPO.java
b/core/src/main/java/org/apache/gravitino/storage/relational/po/TagMetadataObjectRelPO.java
index e26b5cfb76..1a4cde15f3 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/po/TagMetadataObjectRelPO.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/po/TagMetadataObjectRelPO.java
@@ -28,6 +28,7 @@ public class TagMetadataObjectRelPO {
private Long tagId;
private Long metadataObjectId;
private String metadataObjectType;
+ private String tagValue;
private String auditInfo;
private Long currentVersion;
private Long lastVersion;
@@ -49,6 +50,7 @@ public class TagMetadataObjectRelPO {
return Objects.equal(tagId, tagRelPO.tagId)
&& Objects.equal(metadataObjectId, tagRelPO.metadataObjectId)
&& Objects.equal(metadataObjectType, tagRelPO.metadataObjectType)
+ && Objects.equal(tagValue, tagRelPO.tagValue)
&& Objects.equal(auditInfo, tagRelPO.auditInfo)
&& Objects.equal(currentVersion, tagRelPO.currentVersion)
&& Objects.equal(lastVersion, tagRelPO.lastVersion)
@@ -61,6 +63,7 @@ public class TagMetadataObjectRelPO {
tagId,
metadataObjectId,
metadataObjectType,
+ tagValue,
auditInfo,
currentVersion,
lastVersion,
@@ -89,6 +92,11 @@ public class TagMetadataObjectRelPO {
return this;
}
+ public Builder withTagValue(String tagValue) {
+ tagRelPO.tagValue = tagValue;
+ return this;
+ }
+
public Builder withAuditInfo(String auditInfo) {
tagRelPO.auditInfo = auditInfo;
return this;
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/po/TagPO.java
b/core/src/main/java/org/apache/gravitino/storage/relational/po/TagPO.java
index 1bf873fe1a..9ff288cc31 100644
--- a/core/src/main/java/org/apache/gravitino/storage/relational/po/TagPO.java
+++ b/core/src/main/java/org/apache/gravitino/storage/relational/po/TagPO.java
@@ -30,6 +30,8 @@ public class TagPO {
private Long metalakeId;
private String comment;
private String properties;
+ private String allowedValues;
+ private String assignmentValue;
private String auditInfo;
private Long currentVersion;
private Long lastVersion;
@@ -53,6 +55,8 @@ public class TagPO {
&& Objects.equal(metalakeId, tagPO.metalakeId)
&& Objects.equal(comment, tagPO.comment)
&& Objects.equal(properties, tagPO.properties)
+ && Objects.equal(allowedValues, tagPO.allowedValues)
+ && Objects.equal(assignmentValue, tagPO.assignmentValue)
&& Objects.equal(auditInfo, tagPO.auditInfo)
&& Objects.equal(currentVersion, tagPO.currentVersion)
&& Objects.equal(lastVersion, tagPO.lastVersion)
@@ -67,6 +71,8 @@ public class TagPO {
metalakeId,
comment,
properties,
+ allowedValues,
+ assignmentValue,
auditInfo,
currentVersion,
lastVersion,
@@ -105,6 +111,16 @@ public class TagPO {
return this;
}
+ public Builder withAllowedValues(String allowedValues) {
+ tagPO.allowedValues = allowedValues;
+ return this;
+ }
+
+ public Builder withAssignmentValue(String assignmentValue) {
+ tagPO.assignmentValue = assignmentValue;
+ return this;
+ }
+
public Builder withAuditInfo(String auditInfo) {
tagPO.auditInfo = auditInfo;
return this;
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/TagMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/TagMetaService.java
index b56f703d2c..36806449c4 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/TagMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/TagMetaService.java
@@ -20,16 +20,23 @@ package org.apache.gravitino.storage.relational.service;
import static
org.apache.gravitino.metrics.source.MetricsSource.GRAVITINO_RELATIONAL_STORE_METRIC_NAME;
+import com.fasterxml.jackson.core.JsonProcessingException;
import com.google.common.base.Preconditions;
import com.google.common.collect.Lists;
import java.io.IOException;
+import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
+import java.util.Optional;
+import java.util.Set;
import java.util.function.Function;
import java.util.stream.Collectors;
+import javax.annotation.Nullable;
import org.apache.gravitino.Entity;
import org.apache.gravitino.EntityAlreadyExistsException;
import org.apache.gravitino.HasIdentifier;
@@ -38,6 +45,7 @@ import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.Namespace;
import org.apache.gravitino.exceptions.NoSuchEntityException;
import org.apache.gravitino.exceptions.NoSuchTagException;
+import org.apache.gravitino.json.JsonUtils;
import org.apache.gravitino.meta.GenericEntity;
import org.apache.gravitino.meta.TagEntity;
import org.apache.gravitino.metrics.Monitored;
@@ -48,6 +56,8 @@ import org.apache.gravitino.storage.relational.po.TagPO;
import org.apache.gravitino.storage.relational.utils.ExceptionUtils;
import org.apache.gravitino.storage.relational.utils.POConverters;
import org.apache.gravitino.storage.relational.utils.SessionUtils;
+import org.apache.gravitino.tag.TagAssignment;
+import org.apache.gravitino.tag.TagValue;
import org.apache.gravitino.utils.NameIdentifierUtil;
import org.apache.gravitino.utils.NamespaceUtil;
@@ -192,9 +202,7 @@ public class TagMetaService {
throw e;
}
- return tagPOs.stream()
- .map(tagPO -> POConverters.fromTagPO(tagPO,
NamespaceUtil.ofTag(metalake)))
- .collect(Collectors.toList());
+ return tagPOsToTagEntities(tagPOs, NamespaceUtil.ofTag(metalake));
}
@Monitored(
@@ -206,11 +214,11 @@ public class TagMetaService {
MetadataObject metadataObject =
NameIdentifierUtil.toMetadataObject(objectIdent, objectType);
String metalake = objectIdent.namespace().level(0);
- TagPO tagPO = null;
+ List<TagPO> tagPOs = null;
try {
Long metadataObjectId = EntityIdService.getEntityId(objectIdent,
objectType);
- tagPO =
+ tagPOs =
SessionUtils.getWithoutCommit(
TagMetadataObjectRelMapper.class,
mapper ->
@@ -221,14 +229,14 @@ public class TagMetaService {
throw e;
}
- if (tagPO == null) {
+ if (tagPOs == null || tagPOs.isEmpty()) {
throw new NoSuchEntityException(
NoSuchEntityException.NO_SUCH_ENTITY_MESSAGE,
Entity.EntityType.TAG.name().toLowerCase(),
tagIdent.name());
}
- return POConverters.fromTagPO(tagPO, NamespaceUtil.ofTag(metalake));
+ return tagPOsToTagEntities(tagPOs, NamespaceUtil.ofTag(metalake)).get(0);
}
@Monitored(
@@ -236,6 +244,14 @@ public class TagMetaService {
baseMetricName = "listAssociatedMetadataObjectsForTag")
public List<GenericEntity>
listAssociatedMetadataObjectsForTag(NameIdentifier tagIdent)
throws IOException {
+ return listAssociatedMetadataObjectsForTag(tagIdent, null);
+ }
+
+ @Monitored(
+ metricsSource = GRAVITINO_RELATIONAL_STORE_METRIC_NAME,
+ baseMetricName = "listAssociatedMetadataObjectsForTagByValue")
+ public List<GenericEntity> listAssociatedMetadataObjectsForTag(
+ NameIdentifier tagIdent, @Nullable String value) throws IOException {
String metalakeName = tagIdent.namespace().level(0);
String tagName = tagIdent.name();
@@ -244,7 +260,10 @@ public class TagMetaService {
SessionUtils.doWithCommitAndFetchResult(
TagMetadataObjectRelMapper.class,
mapper ->
-
mapper.listTagMetadataObjectRelsByMetalakeAndTagName(metalakeName, tagName));
+ value == null
+ ?
mapper.listTagMetadataObjectRelsByMetalakeAndTagName(metalakeName, tagName)
+ :
mapper.listTagMetadataObjectRelsByMetalakeAndTagNameAndValue(
+ metalakeName, tagName, value));
List<GenericEntity> metadataObjects = Lists.newArrayList();
Map<String, List<TagMetadataObjectRelPO>> tagMetadataObjectRelPOsByType =
@@ -297,64 +316,167 @@ public class TagMetaService {
NameIdentifier[] tagsToAdd,
NameIdentifier[] tagsToRemove)
throws NoSuchEntityException, EntityAlreadyExistsException, IOException {
+ return associateTagValuesWithMetadataObject(
+ objectIdent,
+ objectType,
+ toValuelessTagValues(tagsToAdd),
+ toValuelessTagValues(tagsToRemove),
+ true /* failOnDuplicateValuelessAssignment */);
+ }
+
+ @Monitored(
+ metricsSource = GRAVITINO_RELATIONAL_STORE_METRIC_NAME,
+ baseMetricName = "associateTagValuesWithMetadataObject")
+ public List<TagEntity> associateTagValuesWithMetadataObject(
+ NameIdentifier objectIdent,
+ Entity.EntityType objectType,
+ TagValue[] tagsToAdd,
+ TagValue[] tagsToRemove)
+ throws NoSuchEntityException, EntityAlreadyExistsException, IOException {
+ return associateTagValuesWithMetadataObject(
+ objectIdent,
+ objectType,
+ tagsToAdd,
+ tagsToRemove,
+ false /* failOnDuplicateValuelessAssignment */);
+ }
+
+ private List<TagEntity> associateTagValuesWithMetadataObject(
+ NameIdentifier objectIdent,
+ Entity.EntityType objectType,
+ TagValue[] tagsToAdd,
+ TagValue[] tagsToRemove,
+ boolean failOnDuplicateValuelessAssignment)
+ throws NoSuchEntityException, EntityAlreadyExistsException, IOException {
MetadataObject metadataObject =
NameIdentifierUtil.toMetadataObject(objectIdent, objectType);
String metalake = objectIdent.namespace().level(0);
try {
Long metadataObjectId = EntityIdService.getEntityId(objectIdent,
objectType);
-
- // Fetch all the tags need to associate with the metadata object.
- List<String> tagNamesToAdd =
-
Arrays.stream(tagsToAdd).map(NameIdentifier::name).collect(Collectors.toList());
- List<TagPO> tagPOsToAdd =
- tagNamesToAdd.isEmpty()
+ List<TagValue> tagValuesToAdd = new
ArrayList<>(Arrays.asList(nullToEmpty(tagsToAdd)));
+ List<TagValue> tagValuesToRemove = new
ArrayList<>(Arrays.asList(nullToEmpty(tagsToRemove)));
+ Set<TagValue> commonTagValues = new LinkedHashSet<>(tagValuesToAdd);
+ commonTagValues.retainAll(tagValuesToRemove);
+ tagValuesToAdd.removeAll(commonTagValues);
+ tagValuesToRemove.removeAll(commonTagValues);
+
+ List<String> tagNamesToUpdate = tagNamesToUpdate(tagValuesToAdd,
tagValuesToRemove);
+ List<TagPO> tagPOsToUpdate =
+ tagNamesToUpdate.isEmpty()
? Collections.emptyList()
- : getTagPOsByMetalakeAndNames(metalake, tagNamesToAdd);
+ : getTagPOsByMetalakeAndNames(metalake, tagNamesToUpdate);
+ Map<String, TagPO> tagPOsByName = tagPOsByName(tagPOsToUpdate);
- // Fetch all the tags need to remove from the metadata object.
- List<String> tagNamesToRemove =
-
Arrays.stream(tagsToRemove).map(NameIdentifier::name).collect(Collectors.toList());
- List<TagPO> tagPOsToRemove =
- tagNamesToRemove.isEmpty()
- ? Collections.emptyList()
- : getTagPOsByMetalakeAndNames(metalake, tagNamesToRemove);
+ List<TagPO> currentTagPOs =
+ SessionUtils.getWithoutCommit(
+ TagMetadataObjectRelMapper.class,
+ mapper ->
+ mapper.listTagPOsByMetadataObjectIdAndType(
+ metadataObjectId, metadataObject.type().toString()));
+ Map<Long, Set<Optional<String>>> activeValuesByTagId = new
LinkedHashMap<>();
+ for (TagPO currentTagPO : currentTagPOs) {
+ trackExistingAssignment(currentTagPO, activeValuesByTagId);
+ }
+
+ List<Long> tagIdsToRemove = new ArrayList<>();
+ List<TagMetadataObjectRelPO> tagRelsToRemove = new ArrayList<>();
+ for (TagValue tagValueToRemove : tagValuesToRemove) {
+ TagPO tagPO = tagPOsByName.get(tagValueToRemove.name());
+ if (tagPO == null) {
+ continue;
+ }
+
+ if (tagValueToRemove.value().isPresent()) {
+ activeValuesByTagId
+ .computeIfAbsent(tagPO.getTagId(), ignored -> new
LinkedHashSet<>())
+ .remove(tagValueToRemove.value());
+ tagRelsToRemove.add(
+ tagRelForValue(tagPO, metadataObjectId, metadataObject,
tagValueToRemove));
+ } else {
+ activeValuesByTagId.remove(tagPO.getTagId());
+ tagIdsToRemove.add(tagPO.getTagId());
+ }
+ }
+
+ List<TagMetadataObjectRelPO> tagRelsToAdd = new ArrayList<>();
+ for (TagValue tagValueToAdd : tagValuesToAdd) {
+ TagPO tagPO = tagPOsByName.get(tagValueToAdd.name());
+ if (tagPO == null) {
+ continue;
+ }
+
+ validateAllowedValue(tagPO, tagValueToAdd);
+ Set<Optional<String>> activeValues =
+ activeValuesByTagId.computeIfAbsent(tagPO.getTagId(), ignored ->
new LinkedHashSet<>());
+ Optional<String> value = tagValueToAdd.value();
+ if (activeValues.contains(value)) {
+ if (failOnDuplicateValuelessAssignment && !value.isPresent()) {
+ throw new EntityAlreadyExistsException(
+ "Tag %s is already associated to metadata object %s",
+ tagValueToAdd.name(), metadataObject);
+ }
+ continue;
+ }
+
+ if (value.isPresent()) {
+ if (activeValues.remove(Optional.empty())) {
+ tagRelsToRemove.add(
+ tagRelForValue(
+ tagPO,
+ metadataObjectId,
+ metadataObject,
+ TagValue.noValue(tagValueToAdd.name())));
+ }
+ } else {
+ Preconditions.checkArgument(
+ activeValues.isEmpty(),
+ "Cannot add a valueless assignment for tag %s on metadata object
%s because valued assignments exist",
+ tagValueToAdd.name(),
+ metadataObject);
+ }
+
+ activeValues.add(value);
+ tagRelsToAdd.add(
+ POConverters.initializeTagMetadataObjectRelPOWithVersion(
+ tagPO.getTagId(),
+ metadataObjectId,
+ metadataObject.type().toString(),
+ tagValueToAdd.value().orElse(null)));
+ }
SessionUtils.doMultipleWithCommit(
() -> {
- // Insert the tag metadata object relations.
- if (tagPOsToAdd.isEmpty()) {
+ if (tagIdsToRemove.isEmpty()) {
return;
}
- List<TagMetadataObjectRelPO> tagRelsToAdd =
- tagPOsToAdd.stream()
- .map(
- tagPO ->
-
POConverters.initializeTagMetadataObjectRelPOWithVersion(
- tagPO.getTagId(),
- metadataObjectId,
- metadataObject.type().toString()))
- .collect(Collectors.toList());
SessionUtils.doWithoutCommit(
TagMetadataObjectRelMapper.class,
- mapper ->
mapper.batchInsertTagMetadataObjectRels(tagRelsToAdd));
+ mapper ->
+
mapper.batchDeleteTagMetadataObjectRelsByTagIdsAndMetadataObject(
+ metadataObjectId, metadataObject.type().toString(),
tagIdsToRemove));
},
() -> {
- // Remove the tag metadata object relations.
- if (tagPOsToRemove.isEmpty()) {
+ if (tagRelsToRemove.isEmpty()) {
return;
}
- List<Long> tagIdsToRemove =
-
tagPOsToRemove.stream().map(TagPO::getTagId).collect(Collectors.toList());
SessionUtils.doWithoutCommit(
TagMetadataObjectRelMapper.class,
mapper ->
-
mapper.batchDeleteTagMetadataObjectRelsByTagIdsAndMetadataObject(
- metadataObjectId, metadataObject.type().toString(),
tagIdsToRemove));
+
mapper.batchDeleteTagMetadataObjectRelsByTagIdsAndValuesAndMetadataObject(
+ metadataObjectId, metadataObject.type().toString(),
tagRelsToRemove));
+ },
+ () -> {
+ if (tagRelsToAdd.isEmpty()) {
+ return;
+ }
+
+ SessionUtils.doWithoutCommit(
+ TagMetadataObjectRelMapper.class,
+ mapper ->
mapper.batchInsertTagMetadataObjectRels(tagRelsToAdd));
});
- // Fetch all the tags associated with the metadata object after the
operation.
List<TagPO> tagPOs =
SessionUtils.getWithoutCommit(
TagMetadataObjectRelMapper.class,
@@ -362,9 +484,7 @@ public class TagMetaService {
mapper.listTagPOsByMetadataObjectIdAndType(
metadataObjectId, metadataObject.type().toString()));
- return tagPOs.stream()
- .map(tagPO -> POConverters.fromTagPO(tagPO,
NamespaceUtil.ofTag(metalake)))
- .collect(Collectors.toList());
+ return tagPOsToTagEntities(tagPOs, NamespaceUtil.ofTag(metalake));
} catch (RuntimeException e) {
ExceptionUtils.checkSQLException(e, Entity.EntityType.TAG,
objectIdent.toString());
@@ -394,6 +514,104 @@ public class TagMetaService {
return tagDeletedCount[0] + tagMetadataObjectRelDeletedCount[0];
}
+ private static List<TagEntity> tagPOsToTagEntities(List<TagPO> tagPOs,
Namespace namespace) {
+ Map<Long, List<TagPO>> tagPOsByTagId = new LinkedHashMap<>();
+ for (TagPO tagPO : tagPOs) {
+ tagPOsByTagId.computeIfAbsent(tagPO.getTagId(), ignored -> new
ArrayList<>()).add(tagPO);
+ }
+
+ return tagPOsByTagId.values().stream()
+ .map(tagPOGroup -> tagPOsToTagEntity(tagPOGroup, namespace))
+ .collect(Collectors.toList());
+ }
+
+ private static TagEntity tagPOsToTagEntity(List<TagPO> tagPOGroup, Namespace
namespace) {
+ TagPO firstTagPO = tagPOGroup.get(0);
+ List<String> assignmentValues =
+ tagPOGroup.stream()
+ .map(TagPO::getAssignmentValue)
+ .filter(Objects::nonNull)
+ .distinct()
+ .collect(Collectors.toList());
+
+ TagAssignment assignment =
+ assignmentValues.isEmpty()
+ ? TagAssignment.noValue()
+ : TagAssignment.ofValues(assignmentValues.toArray(new String[0]));
+ return POConverters.fromTagPO(firstTagPO,
namespace).copyWithAssignment(assignment);
+ }
+
+ private static TagValue[] toValuelessTagValues(NameIdentifier[] tags) {
+ if (tags == null) {
+ return null;
+ }
+
+ return Arrays.stream(tags).map(tag ->
TagValue.noValue(tag.name())).toArray(TagValue[]::new);
+ }
+
+ private static TagValue[] nullToEmpty(TagValue[] tagValues) {
+ return tagValues == null ? new TagValue[0] : tagValues;
+ }
+
+ private static List<String> tagNamesToUpdate(
+ List<TagValue> tagsToAdd, List<TagValue> tagsToRemove) {
+ Set<String> tagNames = new LinkedHashSet<>();
+ tagsToAdd.stream().map(TagValue::name).forEach(tagNames::add);
+ tagsToRemove.stream().map(TagValue::name).forEach(tagNames::add);
+ return new ArrayList<>(tagNames);
+ }
+
+ private static Map<String, TagPO> tagPOsByName(List<TagPO> tagPOs) {
+ Map<String, TagPO> tagPOsByName = new LinkedHashMap<>();
+ for (TagPO tagPO : tagPOs) {
+ tagPOsByName.put(tagPO.getTagName(), tagPO);
+ }
+ return tagPOsByName;
+ }
+
+ private static void trackExistingAssignment(
+ TagPO tagPO, Map<Long, Set<Optional<String>>> activeValuesByTagId) {
+ activeValuesByTagId
+ .computeIfAbsent(tagPO.getTagId(), ignored -> new LinkedHashSet<>())
+ .add(Optional.ofNullable(tagPO.getAssignmentValue()));
+ }
+
+ private static TagMetadataObjectRelPO tagRelForValue(
+ TagPO tagPO, Long metadataObjectId, MetadataObject metadataObject,
TagValue tagValue) {
+ return POConverters.initializeTagMetadataObjectRelPOWithVersion(
+ tagPO.getTagId(),
+ metadataObjectId,
+ metadataObject.type().toString(),
+ tagValue.value().orElse(null));
+ }
+
+ private static void validateAllowedValue(TagPO tagPO, TagValue tagValue)
+ throws JsonProcessingException {
+ if (tagPO.getAllowedValues() == null) {
+ return;
+ }
+
+ String[] allowedValues =
+ JsonUtils.anyFieldMapper().readValue(tagPO.getAllowedValues(),
String[].class);
+ if (!tagValue.value().isPresent()) {
+ Preconditions.checkArgument(
+ allowedValues.length == 0,
+ "Tag %s requires assignment values from allowed values %s",
+ tagValue.name(),
+ Arrays.toString(allowedValues));
+ return;
+ }
+
+ Preconditions.checkArgument(
+ allowedValues.length > 0, "Tag %s does not allow assignment values",
tagValue.name());
+ Preconditions.checkArgument(
+ Arrays.asList(allowedValues).contains(tagValue.value().get()),
+ "Tag %s value %s is not in allowed values %s",
+ tagValue.name(),
+ tagValue.value().get(),
+ Arrays.toString(allowedValues));
+ }
+
private TagPO getTagPOByMetalakeAndName(String metalakeName, String tagName)
{
TagPO tagPO =
SessionUtils.getWithoutCommit(
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/utils/POConverters.java
b/core/src/main/java/org/apache/gravitino/storage/relational/utils/POConverters.java
index fcb05a6f46..d019800d1a 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/utils/POConverters.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/utils/POConverters.java
@@ -93,6 +93,7 @@ import org.apache.gravitino.storage.relational.po.TagPO;
import org.apache.gravitino.storage.relational.po.TopicPO;
import org.apache.gravitino.storage.relational.po.UserPO;
import org.apache.gravitino.storage.relational.po.UserRoleRelPO;
+import org.apache.gravitino.tag.TagValueConstraint;
import org.apache.gravitino.utils.PrincipalUtils;
/** POConverters is a utility class to convert PO to Base and vice versa. */
@@ -1413,6 +1414,10 @@ public class POConverters {
.withName(tagPO.getTagName())
.withNamespace(namespace)
.withComment(tagPO.getComment())
+ .withAllowedValues(
+ tagPO.getAllowedValues() == null
+ ? null
+ :
JsonUtils.anyFieldMapper().readValue(tagPO.getAllowedValues(), String[].class))
.withProperties(JsonUtils.anyFieldMapper().readValue(tagPO.getProperties(),
Map.class))
.withAuditInfo(
JsonUtils.anyFieldMapper().readValue(tagPO.getAuditInfo(),
AuditInfo.class))
@@ -1433,6 +1438,7 @@ public class POConverters {
.withTagName(tagEntity.name())
.withComment(tagEntity.comment())
.withProperties(JsonUtils.anyFieldMapper().writeValueAsString(tagEntity.properties()))
+
.withAllowedValues(serializeAllowedValues(tagEntity.valueConstraint()))
.withAuditInfo(JsonUtils.anyFieldMapper().writeValueAsString(tagEntity.auditInfo()))
.withCurrentVersion(INIT_VERSION)
.withLastVersion(INIT_VERSION)
@@ -1455,6 +1461,7 @@ public class POConverters {
.withMetalakeId(oldTagPO.getMetalakeId())
.withComment(newEntity.comment())
.withProperties(JsonUtils.anyFieldMapper().writeValueAsString(newEntity.properties()))
+
.withAllowedValues(serializeAllowedValues(newEntity.valueConstraint()))
.withAuditInfo(JsonUtils.anyFieldMapper().writeValueAsString(newEntity.auditInfo()))
.withCurrentVersion(nextVersion)
.withLastVersion(nextVersion)
@@ -1465,8 +1472,35 @@ public class POConverters {
}
}
+ private static String serializeAllowedValues(TagValueConstraint
valueConstraint)
+ throws JsonProcessingException {
+ String[] allowedValues = allowedValuesForStorage(valueConstraint);
+ return allowedValues == null
+ ? null
+ : JsonUtils.anyFieldMapper().writeValueAsString(allowedValues);
+ }
+
+ private static String[] allowedValuesForStorage(TagValueConstraint
valueConstraint) {
+ // Storage encoding: null means ANY_VALUE, an empty array means NO_VALUE.
+ switch (valueConstraint.type()) {
+ case ANY_VALUE:
+ return null;
+ case NO_VALUE:
+ case ALLOWED_VALUES:
+ return valueConstraint.allowedValues();
+ default:
+ throw new IllegalArgumentException("Unknown tag value constraint: " +
valueConstraint);
+ }
+ }
+
public static TagMetadataObjectRelPO
initializeTagMetadataObjectRelPOWithVersion(
Long tagId, Long metadataObjectId, String metadataObjectType) {
+ return initializeTagMetadataObjectRelPOWithVersion(
+ tagId, metadataObjectId, metadataObjectType, null);
+ }
+
+ public static TagMetadataObjectRelPO
initializeTagMetadataObjectRelPOWithVersion(
+ Long tagId, Long metadataObjectId, String metadataObjectType, String
tagValue) {
try {
AuditInfo auditInfo =
AuditInfo.builder()
@@ -1478,6 +1512,7 @@ public class POConverters {
.withTagId(tagId)
.withMetadataObjectId(metadataObjectId)
.withMetadataObjectType(metadataObjectType)
+ .withTagValue(tagValue == null ? "" : tagValue)
.withAuditInfo(JsonUtils.anyFieldMapper().writeValueAsString(auditInfo))
.withCurrentVersion(INIT_VERSION)
.withLastVersion(INIT_VERSION)
diff --git a/core/src/main/java/org/apache/gravitino/tag/TagManager.java
b/core/src/main/java/org/apache/gravitino/tag/TagManager.java
index c334a338c8..b99721eae1 100644
--- a/core/src/main/java/org/apache/gravitino/tag/TagManager.java
+++ b/core/src/main/java/org/apache/gravitino/tag/TagManager.java
@@ -290,7 +290,7 @@ public class TagManager implements TagDispatcher {
checkMetalake(NameIdentifier.of(metalake), entityStore);
return entityStore
.relationOperations()
- .getEntityByRelation(
+ .<TagEntity>getEntityByRelation(
SupportsRelationOperations.Type.TAG_METADATA_OBJECT_REL,
entityIdent,
entityType,
diff --git a/core/src/test/java/org/apache/gravitino/meta/TestTagEntity.java
b/core/src/test/java/org/apache/gravitino/meta/TestTagEntity.java
new file mode 100644
index 0000000000..9f305f3eda
--- /dev/null
+++ b/core/src/test/java/org/apache/gravitino/meta/TestTagEntity.java
@@ -0,0 +1,63 @@
+/*
+ * 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.meta;
+
+import org.apache.gravitino.Namespace;
+import org.apache.gravitino.tag.TagAssignment;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+public class TestTagEntity {
+
+ @Test
+ public void testTagEntityCarriesAssignmentContextOutsideFields() {
+ TagEntity tagEntity =
+ TagEntity.builder()
+ .withId(1L)
+ .withName("tag")
+ .withNamespace(Namespace.of("metalake"))
+ .withComment("comment")
+ .withAuditInfo(AuditInfo.EMPTY)
+ .withAssignment(TagAssignment.ofValues("dev", "prod"))
+ .build();
+
+ Assertions.assertTrue(tagEntity.assignment().isPresent());
+ Assertions.assertArrayEquals(
+ new String[] {"dev", "prod"}, tagEntity.assignment().get().values());
+
Assertions.assertFalse(tagEntity.fields().values().contains(tagEntity.assignment().get()));
+ }
+
+ @Test
+ public void testTagEntityEqualityIgnoresAssignmentContext() {
+ TagEntity tagEntity =
+ TagEntity.builder()
+ .withId(1L)
+ .withName("tag")
+ .withNamespace(Namespace.of("metalake"))
+ .withComment("comment")
+ .withAuditInfo(AuditInfo.EMPTY)
+ .build();
+
+ TagEntity tagEntityWithAssignment =
+ tagEntity.copyWithAssignment(TagAssignment.ofValues("dev", "prod"));
+
+ Assertions.assertEquals(tagEntity, tagEntityWithAssignment);
+ Assertions.assertEquals(tagEntity.hashCode(),
tagEntityWithAssignment.hashCode());
+ }
+}
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTagMetaService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTagMetaService.java
index 0add50e617..f263a68e12 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTagMetaService.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTagMetaService.java
@@ -30,12 +30,18 @@ import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Statement;
import java.time.Instant;
+import java.util.Arrays;
+import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import org.apache.gravitino.Entity;
import org.apache.gravitino.EntityAlreadyExistsException;
import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.Namespace;
+import org.apache.gravitino.RelationEdgeTarget;
+import org.apache.gravitino.RelationQuery;
+import org.apache.gravitino.RelationUpdate;
+import org.apache.gravitino.SupportsRelationOperations;
import org.apache.gravitino.exceptions.NoSuchEntityException;
import org.apache.gravitino.meta.BaseMetalake;
import org.apache.gravitino.meta.CatalogEntity;
@@ -51,6 +57,7 @@ import org.apache.gravitino.rel.types.Types;
import org.apache.gravitino.storage.RandomIdGenerator;
import org.apache.gravitino.storage.relational.TestJDBCBackend;
import org.apache.gravitino.storage.relational.session.SqlSessionFactoryHelper;
+import org.apache.gravitino.tag.TagValue;
import org.apache.gravitino.utils.NameIdentifierUtil;
import org.apache.gravitino.utils.NamespaceUtil;
import org.apache.ibatis.session.SqlSession;
@@ -59,7 +66,7 @@ import org.junit.jupiter.api.TestTemplate;
public class TestTagMetaService extends TestJDBCBackend {
- private static final String METALAKE_NAME = "metalake_for_tag_test";
+ private static final String METALAKE_NAME =
"metalake_for_tag_meta_service_test";
private final Map<String, String> props = ImmutableMap.of("k1", "v1");
@@ -454,9 +461,9 @@ public class TestTagMetaService extends TestJDBCBackend {
tagMetaService.associateTagsWithMetadataObject(
catalog.nameIdentifier(), catalog.type(), tagsToAdd, new
NameIdentifier[0]);
Assertions.assertEquals(3, tagEntities.size());
- Assertions.assertTrue(tagEntities.contains(tagEntity1));
- Assertions.assertTrue(tagEntities.contains(tagEntity2));
- Assertions.assertTrue(tagEntities.contains(tagEntity3));
+ Assertions.assertTrue(containsValuelessTagAssignment(tagEntities,
tagEntity1));
+ Assertions.assertTrue(containsValuelessTagAssignment(tagEntities,
tagEntity2));
+ Assertions.assertTrue(containsValuelessTagAssignment(tagEntities,
tagEntity3));
// Test disassociate tags with metadata object
NameIdentifier[] tagsToRemove =
@@ -467,18 +474,18 @@ public class TestTagMetaService extends TestJDBCBackend {
catalog.nameIdentifier(), catalog.type(), new NameIdentifier[0],
tagsToRemove);
Assertions.assertEquals(2, tagEntities1.size());
- Assertions.assertFalse(tagEntities1.contains(tagEntity1));
- Assertions.assertTrue(tagEntities1.contains(tagEntity2));
- Assertions.assertTrue(tagEntities1.contains(tagEntity3));
+ Assertions.assertFalse(containsValuelessTagAssignment(tagEntities1,
tagEntity1));
+ Assertions.assertTrue(containsValuelessTagAssignment(tagEntities1,
tagEntity2));
+ Assertions.assertTrue(containsValuelessTagAssignment(tagEntities1,
tagEntity3));
// Test no tags to associate and disassociate
List<TagEntity> tagEntities2 =
tagMetaService.associateTagsWithMetadataObject(
catalog.nameIdentifier(), catalog.type(), new NameIdentifier[0],
new NameIdentifier[0]);
Assertions.assertEquals(2, tagEntities2.size());
- Assertions.assertFalse(tagEntities2.contains(tagEntity1));
- Assertions.assertTrue(tagEntities2.contains(tagEntity2));
- Assertions.assertTrue(tagEntities2.contains(tagEntity3));
+ Assertions.assertFalse(containsValuelessTagAssignment(tagEntities2,
tagEntity1));
+ Assertions.assertTrue(containsValuelessTagAssignment(tagEntities2,
tagEntity2));
+ Assertions.assertTrue(containsValuelessTagAssignment(tagEntities2,
tagEntity3));
// Test associate and disassociate same tags with metadata object
List<TagEntity> tagEntities3 =
@@ -486,9 +493,9 @@ public class TestTagMetaService extends TestJDBCBackend {
catalog.nameIdentifier(), catalog.type(), tagsToRemove,
tagsToRemove);
Assertions.assertEquals(2, tagEntities3.size());
- Assertions.assertFalse(tagEntities3.contains(tagEntity1));
- Assertions.assertTrue(tagEntities3.contains(tagEntity2));
- Assertions.assertTrue(tagEntities3.contains(tagEntity3));
+ Assertions.assertFalse(containsValuelessTagAssignment(tagEntities3,
tagEntity1));
+ Assertions.assertTrue(containsValuelessTagAssignment(tagEntities3,
tagEntity2));
+ Assertions.assertTrue(containsValuelessTagAssignment(tagEntities3,
tagEntity3));
// Test associate and disassociate in-existent tags with metadata object
NameIdentifier[] tagsToAdd1 =
@@ -508,8 +515,8 @@ public class TestTagMetaService extends TestJDBCBackend {
catalog.nameIdentifier(), catalog.type(), tagsToAdd1,
tagsToRemove1);
Assertions.assertEquals(2, tagEntities4.size());
- Assertions.assertTrue(tagEntities4.contains(tagEntity2));
- Assertions.assertTrue(tagEntities4.contains(tagEntity3));
+ Assertions.assertTrue(containsValuelessTagAssignment(tagEntities4,
tagEntity2));
+ Assertions.assertTrue(containsValuelessTagAssignment(tagEntities4,
tagEntity3));
// Test associate already associated tags with metadata object
Assertions.assertThrows(
@@ -524,8 +531,8 @@ public class TestTagMetaService extends TestJDBCBackend {
catalog.nameIdentifier(), catalog.type(), new NameIdentifier[0],
tagsToRemove);
Assertions.assertEquals(2, tagEntities5.size());
- Assertions.assertTrue(tagEntities5.contains(tagEntity2));
- Assertions.assertTrue(tagEntities5.contains(tagEntity3));
+ Assertions.assertTrue(containsValuelessTagAssignment(tagEntities5,
tagEntity2));
+ Assertions.assertTrue(containsValuelessTagAssignment(tagEntities5,
tagEntity3));
// Test associate and disassociate with invalid metadata object
Assertions.assertThrows(
@@ -543,8 +550,8 @@ public class TestTagMetaService extends TestJDBCBackend {
schema.nameIdentifier(), schema.type(), tagsToAdd, tagsToRemove);
Assertions.assertEquals(2, tagEntities6.size());
- Assertions.assertTrue(tagEntities6.contains(tagEntity2));
- Assertions.assertTrue(tagEntities6.contains(tagEntity3));
+ Assertions.assertTrue(containsValuelessTagAssignment(tagEntities6,
tagEntity2));
+ Assertions.assertTrue(containsValuelessTagAssignment(tagEntities6,
tagEntity3));
// Test associate and disassociate to a table
List<TagEntity> tagEntities7 =
@@ -552,8 +559,161 @@ public class TestTagMetaService extends TestJDBCBackend {
table.nameIdentifier(), table.type(), tagsToAdd, tagsToRemove);
Assertions.assertEquals(2, tagEntities7.size());
- Assertions.assertTrue(tagEntities7.contains(tagEntity2));
- Assertions.assertTrue(tagEntities7.contains(tagEntity3));
+ Assertions.assertTrue(containsValuelessTagAssignment(tagEntities7,
tagEntity2));
+ Assertions.assertTrue(containsValuelessTagAssignment(tagEntities7,
tagEntity3));
+ }
+
+ @TestTemplate
+ public void testAssociateTagValuesWithMetadataObject() throws IOException {
+ createAndInsertMakeLake(METALAKE_NAME);
+ CatalogEntity catalog = createAndInsertCatalog(METALAKE_NAME,
"catalog_value");
+
+ TagMetaService tagMetaService = TagMetaService.getInstance();
+ TagEntity tagEntity =
+ TagEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName("stage")
+ .withNamespace(NamespaceUtil.ofTag(METALAKE_NAME))
+ .withComment("stage comment")
+ .withProperties(props)
+ .withAllowedValues(new String[] {"dev", "prod"})
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ tagMetaService.insertTag(tagEntity, false);
+
+ TagEntity storedTag =
+
tagMetaService.getTagByIdentifier(NameIdentifierUtil.ofTag(METALAKE_NAME,
"stage"));
+ Assertions.assertArrayEquals(
+ new String[] {"dev", "prod"},
storedTag.valueConstraint().allowedValues());
+ Assertions.assertFalse(storedTag.assignment().isPresent());
+
+ List<TagEntity> tagEntities =
+ backend.updateEntityRelations(
+ RelationUpdate.of(
+ SupportsRelationOperations.Type.TAG_METADATA_OBJECT_REL,
+ catalog.nameIdentifier(),
+ catalog.type(),
+ new RelationEdgeTarget[] {
+ RelationEdgeTarget.of(
+ NameIdentifierUtil.ofTag(METALAKE_NAME, "stage"),
+ Entity.EntityType.TAG,
+ "dev"),
+ RelationEdgeTarget.of(
+ NameIdentifierUtil.ofTag(METALAKE_NAME, "stage"),
+ Entity.EntityType.TAG,
+ "prod")
+ },
+ new RelationEdgeTarget[0]));
+
+ Assertions.assertEquals(1, tagEntities.size());
+ Assertions.assertEquals("stage", tagEntities.get(0).name());
+ assertAssignmentValues(tagEntities.get(0), "dev", "prod");
+ Assertions.assertEquals(2, countActiveTagRel(tagEntity.id()));
+
+ TagEntity tagForMetadataObject =
+ tagMetaService.getTagForMetadataObject(
+ catalog.nameIdentifier(),
+ catalog.type(),
+ NameIdentifierUtil.ofTag(METALAKE_NAME, "stage"));
+ assertAssignmentValues(tagForMetadataObject, "dev", "prod");
+
+ List<GenericEntity> devMetadataObjects =
+ backend.listEntitiesByRelation(
+ RelationQuery.of(
+ SupportsRelationOperations.Type.TAG_METADATA_OBJECT_REL,
+ NameIdentifierUtil.ofTag(METALAKE_NAME, "stage"),
+ Entity.EntityType.TAG,
+ true,
+ "dev"));
+ Assertions.assertEquals(1, devMetadataObjects.size());
+ Assertions.assertTrue(
+ containsGenericEntity(devMetadataObjects, "catalog_value",
Entity.EntityType.CATALOG));
+
+ List<GenericEntity> missingMetadataObjects =
+ backend.listEntitiesByRelation(
+ RelationQuery.of(
+ SupportsRelationOperations.Type.TAG_METADATA_OBJECT_REL,
+ NameIdentifierUtil.ofTag(METALAKE_NAME, "stage"),
+ Entity.EntityType.TAG,
+ true,
+ "missing"));
+ Assertions.assertEquals(0, missingMetadataObjects.size());
+
+ List<TagEntity> tagEntitiesAfterRemove =
+ tagMetaService.associateTagValuesWithMetadataObject(
+ catalog.nameIdentifier(),
+ catalog.type(),
+ new TagValue[0],
+ new TagValue[] {TagValue.of("stage", "dev")});
+ Assertions.assertEquals(1, tagEntitiesAfterRemove.size());
+ assertAssignmentValues(tagEntitiesAfterRemove.get(0), "prod");
+ Assertions.assertEquals(1, countActiveTagRel(tagEntity.id()));
+
+ List<TagEntity> tagEntitiesAfterDuplicateAdd =
+ tagMetaService.associateTagValuesWithMetadataObject(
+ catalog.nameIdentifier(),
+ catalog.type(),
+ new TagValue[] {TagValue.of("stage", "prod")},
+ new TagValue[0]);
+ Assertions.assertEquals(1, tagEntitiesAfterDuplicateAdd.size());
+ assertAssignmentValues(tagEntitiesAfterDuplicateAdd.get(0), "prod");
+ Assertions.assertEquals(1, countActiveTagRel(tagEntity.id()));
+
+ List<TagEntity> tagEntitiesAfterSameValueAddRemove =
+ tagMetaService.associateTagValuesWithMetadataObject(
+ catalog.nameIdentifier(),
+ catalog.type(),
+ new TagValue[] {TagValue.of("stage", "prod")},
+ new TagValue[] {TagValue.of("stage", "prod")});
+ Assertions.assertEquals(1, tagEntitiesAfterSameValueAddRemove.size());
+ assertAssignmentValues(tagEntitiesAfterSameValueAddRemove.get(0), "prod");
+ Assertions.assertEquals(1, countActiveTagRel(tagEntity.id()));
+
+ IllegalArgumentException exception =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ tagMetaService.associateTagValuesWithMetadataObject(
+ catalog.nameIdentifier(),
+ catalog.type(),
+ new TagValue[] {TagValue.of("stage", "qa")},
+ new TagValue[0]));
+ Assertions.assertTrue(exception.getMessage().contains("is not in allowed
values"));
+
+ IllegalArgumentException noValueException =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
+ tagMetaService.associateTagValuesWithMetadataObject(
+ catalog.nameIdentifier(),
+ catalog.type(),
+ new TagValue[] {TagValue.noValue("stage")},
+ new TagValue[0]));
+ Assertions.assertTrue(noValueException.getMessage().contains("requires
assignment values"));
+ }
+
+ @TestTemplate
+ public void testRejectDuplicateValuelessTagAssignment() throws IOException {
+ createAndInsertMakeLake(METALAKE_NAME);
+ CatalogEntity catalog = createAndInsertCatalog(METALAKE_NAME,
"catalog_unique_value");
+
+ TagEntity tagEntity =
+ TagEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName("unique_value")
+ .withNamespace(NamespaceUtil.ofTag(METALAKE_NAME))
+ .withComment("unique value comment")
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ TagMetaService tagMetaService = TagMetaService.getInstance();
+ tagMetaService.insertTag(tagEntity, false);
+ tagMetaService.associateTagValuesWithMetadataObject(
+ catalog.nameIdentifier(),
+ catalog.type(),
+ new TagValue[] {TagValue.noValue(tagEntity.name())},
+ new TagValue[0]);
+
+ Assertions.assertThrows(SQLException.class, () ->
insertDuplicateActiveTagRel(tagEntity.id()));
}
@TestTemplate
@@ -1124,6 +1284,24 @@ public class TestTagMetaService extends TestJDBCBackend {
() -> tagMetaService.getTagIdByTagName(metalakeId, "missing_tag"));
}
+ private boolean containsValuelessTagAssignment(
+ List<TagEntity> tagEntities, TagEntity expectedTagEntity) {
+ return tagEntities.stream()
+ .anyMatch(
+ tagEntity ->
+ tagEntity.id().equals(expectedTagEntity.id())
+ && tagEntity.name().equals(expectedTagEntity.name())
+ && tagEntity.assignment().isPresent()
+ && !tagEntity.assignment().get().hasValues());
+ }
+
+ private void assertAssignmentValues(TagEntity tagEntity, String...
expectedValues) {
+ assertTrue(tagEntity.assignment().isPresent());
+ assertEquals(
+ new LinkedHashSet<>(Arrays.asList(expectedValues)),
+ new
LinkedHashSet<>(Arrays.asList(tagEntity.assignment().get().values())));
+ }
+
private boolean containsGenericEntity(
List<GenericEntity> genericEntities, String name, Entity.EntityType
entityType) {
return genericEntities.stream().anyMatch(e -> e.name().equals(name) &&
e.type() == entityType);
@@ -1147,6 +1325,20 @@ public class TestTagMetaService extends TestJDBCBackend {
}
}
+ private void insertDuplicateActiveTagRel(Long tagId) throws SQLException {
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ Statement statement = connection.createStatement()) {
+ statement.executeUpdate(
+ String.format(
+ "INSERT INTO tag_relation_meta (tag_id, metadata_object_id,
metadata_object_type, tag_value, audit_info, current_version, last_version,
deleted_at) "
+ + "SELECT tag_id, metadata_object_id, metadata_object_type,
tag_value, audit_info, current_version, last_version, deleted_at "
+ + "FROM tag_relation_meta WHERE tag_id = %d AND deleted_at =
0",
+ tagId));
+ }
+ }
+
private Integer countActiveTagRel(Long tagId) {
try (SqlSession sqlSession =
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/utils/TestPOConverters.java
b/core/src/test/java/org/apache/gravitino/storage/relational/utils/TestPOConverters.java
index bff96b2bd6..9ce171d878 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/utils/TestPOConverters.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/utils/TestPOConverters.java
@@ -891,6 +891,7 @@ public class TestPOConverters {
assertEquals(1L, tagMetadataObjectRelPO.getMetadataObjectId());
assertEquals(
MetadataObject.Type.CATALOG.toString(),
tagMetadataObjectRelPO.getMetadataObjectType());
+ assertEquals("", tagMetadataObjectRelPO.getTagValue());
assertEquals(1, tagMetadataObjectRelPO.getCurrentVersion());
assertEquals(1, tagMetadataObjectRelPO.getLastVersion());
diff --git a/scripts/h2/schema-2.0.0-h2.sql b/scripts/h2/schema-2.0.0-h2.sql
index 08ea0bcd20..05039ccd48 100644
--- a/scripts/h2/schema-2.0.0-h2.sql
+++ b/scripts/h2/schema-2.0.0-h2.sql
@@ -295,6 +295,7 @@ CREATE TABLE IF NOT EXISTS `tag_meta` (
`metalake_id` BIGINT(20) UNSIGNED NOT NULL COMMENT 'metalake id',
`tag_comment` VARCHAR(256) DEFAULT '' COMMENT 'tag comment',
`properties` CLOB DEFAULT NULL COMMENT 'tag properties',
+ `allowed_values` CLOB DEFAULT NULL COMMENT 'tag allowed values as a JSON
string array, NULL allows any value, [] allows no value',
`audit_info` CLOB NOT NULL COMMENT 'tag audit info',
`current_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'tag current
version',
`last_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'tag last version',
@@ -308,14 +309,16 @@ CREATE TABLE IF NOT EXISTS `tag_relation_meta` (
`tag_id` BIGINT(20) UNSIGNED NOT NULL COMMENT 'tag id',
`metadata_object_id` BIGINT(20) UNSIGNED NOT NULL COMMENT 'metadata object
id',
`metadata_object_type` VARCHAR(64) NOT NULL COMMENT 'metadata object type',
+ `tag_value` VARCHAR(256) NOT NULL DEFAULT '' COMMENT 'tag assignment
value, empty string means no value',
`audit_info` CLOB NOT NULL COMMENT 'tag relation audit info',
`current_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'tag relation
current version',
`last_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'tag relation last
version',
`deleted_at` BIGINT(20) UNSIGNED NOT NULL DEFAULT 0 COMMENT 'tag relation
deleted at',
PRIMARY KEY (`id`),
- UNIQUE KEY `uk_ti_mi_del` (`tag_id`, `metadata_object_id`, `deleted_at`),
+ UNIQUE KEY `uk_ti_mi_mo_tv_del` (`tag_id`, `metadata_object_id`,
`metadata_object_type`, `tag_value`, `deleted_at`),
KEY `idx_tid` (`tag_id`),
- KEY `idx_mid` (`metadata_object_id`)
+ KEY `idx_mid` (`metadata_object_id`),
+ KEY `idx_tid_value` (`tag_id`, `tag_value`)
) ENGINE=InnoDB;
CREATE TABLE IF NOT EXISTS `owner_meta` (
diff --git a/scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql
b/scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql
index 16ea9b0392..f7a856845e 100644
--- a/scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql
+++ b/scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql
@@ -27,3 +27,12 @@ CREATE UNIQUE INDEX IF NOT EXISTS `uk_mid_geid_del` ON
`group_meta` (`metalake_i
ALTER TABLE `table_column_version_info`
ALTER COLUMN `column_comment` VARCHAR(4096) DEFAULT '';
+
+ALTER TABLE `tag_meta` ADD COLUMN `allowed_values` CLOB DEFAULT NULL COMMENT
'tag allowed values as a JSON string array, NULL allows any value, [] allows no
value' AFTER `properties`;
+
+ALTER TABLE `tag_relation_meta` DROP INDEX `uk_ti_mi_del`;
+
+ALTER TABLE `tag_relation_meta` ADD COLUMN `tag_value` VARCHAR(256) NOT NULL
DEFAULT '' COMMENT 'tag assignment value, empty string means no value' AFTER
`metadata_object_type`;
+
+CREATE UNIQUE INDEX IF NOT EXISTS `uk_ti_mi_mo_tv_del` ON `tag_relation_meta`
(`tag_id`, `metadata_object_id`, `metadata_object_type`, `tag_value`,
`deleted_at`);
+CREATE INDEX IF NOT EXISTS `idx_tid_value` ON `tag_relation_meta` (`tag_id`,
`tag_value`);
diff --git a/scripts/mysql/schema-2.0.0-mysql.sql
b/scripts/mysql/schema-2.0.0-mysql.sql
index 06ac92625b..5defbd4faf 100644
--- a/scripts/mysql/schema-2.0.0-mysql.sql
+++ b/scripts/mysql/schema-2.0.0-mysql.sql
@@ -286,6 +286,7 @@ CREATE TABLE IF NOT EXISTS `tag_meta` (
`metalake_id` BIGINT(20) UNSIGNED NOT NULL COMMENT 'metalake id',
`tag_comment` VARCHAR(256) DEFAULT '' COMMENT 'tag comment',
`properties` MEDIUMTEXT DEFAULT NULL COMMENT 'tag properties',
+ `allowed_values` MEDIUMTEXT DEFAULT NULL COMMENT 'tag allowed values as a
JSON string array, NULL allows any value, [] allows no value',
`audit_info` MEDIUMTEXT NOT NULL COMMENT 'tag audit info',
`current_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'tag current
version',
`last_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'tag last version',
@@ -299,14 +300,16 @@ CREATE TABLE IF NOT EXISTS `tag_relation_meta` (
`tag_id` BIGINT(20) UNSIGNED NOT NULL COMMENT 'tag id',
`metadata_object_id` BIGINT(20) UNSIGNED NOT NULL COMMENT 'metadata object
id',
`metadata_object_type` VARCHAR(64) NOT NULL COMMENT 'metadata object type',
+ `tag_value` VARCHAR(256) NOT NULL DEFAULT '' COMMENT 'tag assignment
value, empty string means no value',
`audit_info` MEDIUMTEXT NOT NULL COMMENT 'tag relation audit info',
`current_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'tag relation
current version',
`last_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'tag relation last
version',
`deleted_at` BIGINT(20) UNSIGNED NOT NULL DEFAULT 0 COMMENT 'tag relation
deleted at',
PRIMARY KEY (`id`),
- UNIQUE KEY `uk_ti_mi_mo_del` (`tag_id`, `metadata_object_id`,
`metadata_object_type`, `deleted_at`),
+ UNIQUE KEY `uk_ti_mi_mo_tv_del` (`tag_id`, `metadata_object_id`,
`metadata_object_type`, `tag_value`, `deleted_at`),
KEY `idx_tid` (`tag_id`),
- KEY `idx_mid` (`metadata_object_id`)
+ KEY `idx_mid` (`metadata_object_id`),
+ KEY `idx_tid_value` (`tag_id`, `tag_value`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_bin COMMENT 'tag
metadata object relation';
CREATE TABLE IF NOT EXISTS `owner_meta` (
diff --git a/scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql
b/scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql
index db34c6e179..32a0f84158 100644
--- a/scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql
+++ b/scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql
@@ -29,3 +29,15 @@ CREATE UNIQUE INDEX `uk_mid_geid_del` ON `group_meta`
(`metalake_id`, `external_
ALTER TABLE `table_column_version_info`
MODIFY COLUMN `column_comment` VARCHAR(4096) DEFAULT '' COMMENT 'column
comment';
+
+ALTER TABLE `tag_meta`
+ ADD COLUMN `allowed_values` MEDIUMTEXT DEFAULT NULL COMMENT 'tag allowed
values as a JSON string array, NULL allows any value, [] allows no value' AFTER
`properties`;
+
+ALTER TABLE `tag_relation_meta`
+ DROP INDEX `uk_ti_mi_mo_del`;
+
+ALTER TABLE `tag_relation_meta`
+ ADD COLUMN `tag_value` VARCHAR(256) NOT NULL DEFAULT '' COMMENT 'tag
assignment value, empty string means no value' AFTER `metadata_object_type`;
+
+CREATE UNIQUE INDEX `uk_ti_mi_mo_tv_del` ON `tag_relation_meta` (`tag_id`,
`metadata_object_id`, `metadata_object_type`, `tag_value`, `deleted_at`);
+CREATE INDEX `idx_tid_value` ON `tag_relation_meta` (`tag_id`, `tag_value`);
diff --git a/scripts/postgresql/schema-2.0.0-postgresql.sql
b/scripts/postgresql/schema-2.0.0-postgresql.sql
index a7d82cc547..cbc05204e8 100644
--- a/scripts/postgresql/schema-2.0.0-postgresql.sql
+++ b/scripts/postgresql/schema-2.0.0-postgresql.sql
@@ -500,6 +500,7 @@ CREATE TABLE IF NOT EXISTS tag_meta (
metalake_id BIGINT NOT NULL,
tag_comment VARCHAR(256) DEFAULT '',
properties TEXT DEFAULT NULL,
+ allowed_values TEXT DEFAULT NULL,
audit_info TEXT NOT NULL,
current_version INT NOT NULL DEFAULT 1,
last_version INT NOT NULL DEFAULT 1,
@@ -515,6 +516,7 @@ COMMENT ON COLUMN tag_meta.tag_name IS 'tag name';
COMMENT ON COLUMN tag_meta.metalake_id IS 'metalake id';
COMMENT ON COLUMN tag_meta.tag_comment IS 'tag comment';
COMMENT ON COLUMN tag_meta.properties IS 'tag properties';
+COMMENT ON COLUMN tag_meta.allowed_values IS 'tag allowed values as a JSON
string array, NULL allows any value, [] allows no value';
COMMENT ON COLUMN tag_meta.audit_info IS 'tag audit info';
@@ -523,21 +525,24 @@ CREATE TABLE IF NOT EXISTS tag_relation_meta (
tag_id BIGINT NOT NULL,
metadata_object_id BIGINT NOT NULL,
metadata_object_type VARCHAR(64) NOT NULL,
+ tag_value VARCHAR(256) NOT NULL DEFAULT '',
audit_info TEXT NOT NULL,
current_version INT NOT NULL DEFAULT 1,
last_version INT NOT NULL DEFAULT 1,
deleted_at BIGINT NOT NULL DEFAULT 0,
PRIMARY KEY (id),
- UNIQUE (tag_id, metadata_object_id, metadata_object_type, deleted_at)
+ CONSTRAINT uk_ti_mi_mo_tv_del UNIQUE (tag_id, metadata_object_id,
metadata_object_type, tag_value, deleted_at)
);
CREATE INDEX IF NOT EXISTS tag_relation_meta_idx_tag_id ON tag_relation_meta
(tag_id);
CREATE INDEX IF NOT EXISTS tag_relation_meta_idx_metadata_object_id ON
tag_relation_meta (metadata_object_id);
+CREATE INDEX IF NOT EXISTS tag_relation_meta_idx_tag_id_value ON
tag_relation_meta (tag_id, tag_value);
COMMENT ON TABLE tag_relation_meta IS 'tag metadata object relation';
COMMENT ON COLUMN tag_relation_meta.id IS 'auto increment id';
COMMENT ON COLUMN tag_relation_meta.tag_id IS 'tag id';
COMMENT ON COLUMN tag_relation_meta.metadata_object_id IS 'metadata object id';
COMMENT ON COLUMN tag_relation_meta.metadata_object_type IS 'metadata object
type';
+COMMENT ON COLUMN tag_relation_meta.tag_value IS 'tag assignment value, empty
string means no value';
COMMENT ON COLUMN tag_relation_meta.audit_info IS 'tag relation audit info';
COMMENT ON COLUMN tag_relation_meta.current_version IS 'tag relation current
version';
COMMENT ON COLUMN tag_relation_meta.last_version IS 'tag relation last
version';
diff --git a/scripts/postgresql/upgrade-1.3.0-to-2.0.0-postgresql.sql
b/scripts/postgresql/upgrade-1.3.0-to-2.0.0-postgresql.sql
index 8c9cd7556c..73c04328a7 100644
--- a/scripts/postgresql/upgrade-1.3.0-to-2.0.0-postgresql.sql
+++ b/scripts/postgresql/upgrade-1.3.0-to-2.0.0-postgresql.sql
@@ -31,3 +31,14 @@ CREATE UNIQUE INDEX IF NOT EXISTS uk_mid_geid_del ON
group_meta (metalake_id, ex
ALTER TABLE table_column_version_info
ALTER COLUMN column_comment TYPE VARCHAR(4096);
+
+ALTER TABLE tag_meta ADD COLUMN IF NOT EXISTS allowed_values TEXT DEFAULT NULL;
+COMMENT ON COLUMN tag_meta.allowed_values IS 'tag allowed values as a JSON
string array, NULL allows any value, [] allows no value';
+
+ALTER TABLE tag_relation_meta ADD COLUMN IF NOT EXISTS tag_value VARCHAR(256)
NOT NULL DEFAULT '';
+COMMENT ON COLUMN tag_relation_meta.tag_value IS 'tag assignment value, empty
string means no value';
+
+ALTER TABLE tag_relation_meta DROP CONSTRAINT IF EXISTS
tag_relation_meta_tag_id_metadata_object_id_metadata_object_key;
+
+CREATE UNIQUE INDEX IF NOT EXISTS uk_ti_mi_mo_tv_del ON tag_relation_meta
(tag_id, metadata_object_id, metadata_object_type, tag_value, deleted_at);
+CREATE INDEX IF NOT EXISTS tag_relation_meta_idx_tag_id_value ON
tag_relation_meta (tag_id, tag_value);