This is an automated email from the ASF dual-hosted git repository. yuqi1129 pushed a commit to branch feat/entity-change-log in repository https://gitbox.apache.org/repos/asf/gravitino.git
commit 2a3abf2ec7e5dce767f070949470326cc2baf43d Author: yuqi <[email protected]> AuthorDate: Fri May 8 16:15:08 2026 +0800 fix(cache): address PR #10914 review feedback - Remove unused OperateType.INSERT and switch fromCode to O(1) map lookup - Rename OperateTypeTypeHandler -> OperateTypeHandler - Move EntityChangeLogPostgreSQLProvider to provider/postgresql/ standalone file - ViewMetaService: separate changelog insert into its own top-level lambda - TopicMetaService/FilesetMetaService: guard changelog insert with deleteResult and return deleteResult > 0, matching TableMetaService pattern - TestEntityChangeLogMapper: use PreparedStatement in forceCreatedAt helper --- .../mapper/EntityChangeLogSQLProviderFactory.java | 27 +----------- ...ypeTypeHandler.java => OperateTypeHandler.java} | 2 +- .../EntityChangeLogPostgreSQLProvider.java | 50 +++++++++++++++++++++ .../storage/relational/po/cache/OperateType.java | 20 ++++++--- .../relational/service/FilesetMetaService.java | 51 ++++++++++++---------- .../relational/service/TopicMetaService.java | 46 ++++++++++--------- .../relational/service/ViewMetaService.java | 4 ++ .../session/SqlSessionFactoryHelper.java | 6 +-- .../provider/base/TestEntityChangeLogMapper.java | 24 +++++----- 9 files changed, 136 insertions(+), 94 deletions(-) diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/EntityChangeLogSQLProviderFactory.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/EntityChangeLogSQLProviderFactory.java index 251db66f98..4ee8e43156 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/EntityChangeLogSQLProviderFactory.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/EntityChangeLogSQLProviderFactory.java @@ -18,12 +18,11 @@ */ package org.apache.gravitino.storage.relational.mapper; -import static org.apache.gravitino.storage.relational.mapper.EntityChangeLogMapper.ENTITY_CHANGE_LOG_TABLE_NAME; - import com.google.common.collect.ImmutableMap; import java.util.Map; import org.apache.gravitino.storage.relational.JDBCBackend.JDBCBackendType; import org.apache.gravitino.storage.relational.mapper.provider.base.EntityChangeLogBaseSQLProvider; +import org.apache.gravitino.storage.relational.mapper.provider.postgresql.EntityChangeLogPostgreSQLProvider; import org.apache.gravitino.storage.relational.po.cache.OperateType; import org.apache.gravitino.storage.relational.session.SqlSessionFactoryHelper; import org.apache.ibatis.annotations.Param; @@ -51,30 +50,6 @@ public class EntityChangeLogSQLProviderFactory { static class EntityChangeLogH2Provider extends EntityChangeLogBaseSQLProvider {} - static class EntityChangeLogPostgreSQLProvider extends EntityChangeLogBaseSQLProvider { - @Override - public String insertEntityChange( - @Param("metalakeName") String metalakeName, - @Param("entityType") String entityType, - @Param("fullName") String fullName, - @Param("operateType") OperateType operateType) { - return "INSERT INTO " - + ENTITY_CHANGE_LOG_TABLE_NAME - + " (metalake_name, entity_type, entity_full_name, operate_type, created_at)" - + " VALUES (#{metalakeName}, #{entityType}, #{fullName}, #{operateType}," - + " CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000 AS BIGINT))"; - } - - @Override - public String pruneOldEntityChanges(@Param("before") long before) { - return "DELETE FROM " - + ENTITY_CHANGE_LOG_TABLE_NAME - + " WHERE id IN (SELECT id FROM " - + ENTITY_CHANGE_LOG_TABLE_NAME - + " WHERE created_at < #{before} ORDER BY created_at LIMIT 1000)"; - } - } - public static String selectEntityChanges( @Param("createdAtFrom") long createdAtFrom, @Param("maxRows") int maxRows) { return getProvider().selectEntityChanges(createdAtFrom, maxRows); diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OperateTypeTypeHandler.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OperateTypeHandler.java similarity index 96% rename from core/src/main/java/org/apache/gravitino/storage/relational/mapper/OperateTypeTypeHandler.java rename to core/src/main/java/org/apache/gravitino/storage/relational/mapper/OperateTypeHandler.java index 1389c4ca62..9e582ebed0 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OperateTypeTypeHandler.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OperateTypeHandler.java @@ -32,7 +32,7 @@ import org.apache.ibatis.type.MappedTypes; * {@code entity_change_log.operate_type} column. */ @MappedTypes(OperateType.class) -public class OperateTypeTypeHandler extends BaseTypeHandler<OperateType> { +public class OperateTypeHandler extends BaseTypeHandler<OperateType> { @Override public void setNonNullParameter( diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/EntityChangeLogPostgreSQLProvider.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/EntityChangeLogPostgreSQLProvider.java new file mode 100644 index 0000000000..d6cdbc5851 --- /dev/null +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/EntityChangeLogPostgreSQLProvider.java @@ -0,0 +1,50 @@ +/* + * 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.storage.relational.mapper.provider.postgresql; + +import static org.apache.gravitino.storage.relational.mapper.EntityChangeLogMapper.ENTITY_CHANGE_LOG_TABLE_NAME; + +import org.apache.gravitino.storage.relational.mapper.provider.base.EntityChangeLogBaseSQLProvider; +import org.apache.gravitino.storage.relational.po.cache.OperateType; +import org.apache.ibatis.annotations.Param; + +public class EntityChangeLogPostgreSQLProvider extends EntityChangeLogBaseSQLProvider { + + @Override + public String insertEntityChange( + @Param("metalakeName") String metalakeName, + @Param("entityType") String entityType, + @Param("fullName") String fullName, + @Param("operateType") OperateType operateType) { + return "INSERT INTO " + + ENTITY_CHANGE_LOG_TABLE_NAME + + " (metalake_name, entity_type, entity_full_name, operate_type, created_at)" + + " VALUES (#{metalakeName}, #{entityType}, #{fullName}, #{operateType}," + + " CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000 AS BIGINT))"; + } + + @Override + public String pruneOldEntityChanges(@Param("before") long before) { + return "DELETE FROM " + + ENTITY_CHANGE_LOG_TABLE_NAME + + " WHERE id IN (SELECT id FROM " + + ENTITY_CHANGE_LOG_TABLE_NAME + + " WHERE created_at < #{before} ORDER BY created_at LIMIT 1000)"; + } +} diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/po/cache/OperateType.java b/core/src/main/java/org/apache/gravitino/storage/relational/po/cache/OperateType.java index 2a2122563d..943ad80d5c 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/po/cache/OperateType.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/po/cache/OperateType.java @@ -18,6 +18,11 @@ */ package org.apache.gravitino.storage.relational.po.cache; +import java.util.Arrays; +import java.util.Map; +import java.util.function.Function; +import java.util.stream.Collectors; + /** * Operate type emitted into {@code entity_change_log.operate_type}. * @@ -28,8 +33,10 @@ package org.apache.gravitino.storage.relational.po.cache; */ public enum OperateType { ALTER(1), - DROP(2), - INSERT(3); + DROP(2); + + private static final Map<Integer, OperateType> BY_CODE = + Arrays.stream(values()).collect(Collectors.toMap(OperateType::getCode, Function.identity())); private final int code; @@ -42,11 +49,10 @@ public enum OperateType { } public static OperateType fromCode(int code) { - for (OperateType type : values()) { - if (type.code == code) { - return type; - } + OperateType type = BY_CODE.get(code); + if (type == null) { + throw new IllegalArgumentException("Unknown OperateType code: " + code); } - throw new IllegalArgumentException("Unknown OperateType code: " + code); + return type; } } diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/service/FilesetMetaService.java b/core/src/main/java/org/apache/gravitino/storage/relational/service/FilesetMetaService.java index 51291a387d..9c1f677b4b 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/service/FilesetMetaService.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/service/FilesetMetaService.java @@ -24,6 +24,7 @@ import com.google.common.base.Preconditions; import java.io.IOException; import java.util.List; import java.util.Objects; +import java.util.concurrent.atomic.AtomicInteger; import java.util.function.Function; import java.util.stream.Collectors; import org.apache.gravitino.Entity; @@ -319,55 +320,57 @@ public class FilesetMetaService { .toString(); // We should delete meta and version info + AtomicInteger deleteResult = new AtomicInteger(0); SessionUtils.doMultipleWithCommit( () -> - SessionUtils.doWithoutCommit( - FilesetMetaMapper.class, - mapper -> mapper.softDeleteFilesetMetasByFilesetId(filesetId)), - () -> + deleteResult.set( + SessionUtils.getWithoutCommit( + FilesetMetaMapper.class, + mapper -> mapper.softDeleteFilesetMetasByFilesetId(filesetId))), + () -> { + if (deleteResult.get() > 0) { SessionUtils.doWithoutCommit( FilesetVersionMapper.class, - mapper -> mapper.softDeleteFilesetVersionsByFilesetId(filesetId)), - () -> + mapper -> mapper.softDeleteFilesetVersionsByFilesetId(filesetId)); SessionUtils.doWithoutCommit( OwnerMetaMapper.class, mapper -> mapper.softDeleteOwnerRelByMetadataObjectIdAndType( - filesetId, MetadataObject.Type.FILESET.name())), - () -> + filesetId, MetadataObject.Type.FILESET.name())); SessionUtils.doWithoutCommit( SecurableObjectMapper.class, mapper -> mapper.softDeleteObjectRelsByMetadataObject( - filesetId, MetadataObject.Type.FILESET.name())), - () -> + filesetId, MetadataObject.Type.FILESET.name())); SessionUtils.doWithoutCommit( TagMetadataObjectRelMapper.class, mapper -> mapper.softDeleteTagMetadataObjectRelsByMetadataObject( - filesetId, MetadataObject.Type.FILESET.name())), - () -> + filesetId, MetadataObject.Type.FILESET.name())); SessionUtils.doWithoutCommit( StatisticMetaMapper.class, - mapper -> mapper.softDeleteStatisticsByEntityId(filesetId)), - () -> + mapper -> mapper.softDeleteStatisticsByEntityId(filesetId)); SessionUtils.doWithoutCommit( PolicyMetadataObjectRelMapper.class, mapper -> mapper.softDeletePolicyMetadataObjectRelsByMetadataObject( - filesetId, MetadataObject.Type.FILESET.name())), + filesetId, MetadataObject.Type.FILESET.name())); + } + }, () -> { - SessionUtils.doWithoutCommit( - EntityChangeLogMapper.class, - mapper -> - mapper.insertEntityChange( - metalakeName, - Entity.EntityType.FILESET.name(), - filesetFullName, - OperateType.DROP)); + if (deleteResult.get() > 0) { + SessionUtils.doWithoutCommit( + EntityChangeLogMapper.class, + mapper -> + mapper.insertEntityChange( + metalakeName, + Entity.EntityType.FILESET.name(), + filesetFullName, + OperateType.DROP)); + } }); - return true; + return deleteResult.get() > 0; } @Monitored( diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/service/TopicMetaService.java b/core/src/main/java/org/apache/gravitino/storage/relational/service/TopicMetaService.java index fd0fbe60ea..85bc2dbbe0 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/service/TopicMetaService.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/service/TopicMetaService.java @@ -300,50 +300,54 @@ public class TopicMetaService { NameIdentifierUtil.ofTopic(metalakeName, catalogName, schemaName, identifier.name()) .toString(); + AtomicInteger deleteResult = new AtomicInteger(0); SessionUtils.doMultipleWithCommit( () -> - SessionUtils.doWithoutCommit( - TopicMetaMapper.class, mapper -> mapper.softDeleteTopicMetasByTopicId(topicId)), - () -> + deleteResult.set( + SessionUtils.getWithoutCommit( + TopicMetaMapper.class, + mapper -> mapper.softDeleteTopicMetasByTopicId(topicId))), + () -> { + if (deleteResult.get() > 0) { SessionUtils.doWithoutCommit( OwnerMetaMapper.class, mapper -> mapper.softDeleteOwnerRelByMetadataObjectIdAndType( - topicId, MetadataObject.Type.TOPIC.name())), - () -> + topicId, MetadataObject.Type.TOPIC.name())); SessionUtils.doWithoutCommit( SecurableObjectMapper.class, mapper -> mapper.softDeleteObjectRelsByMetadataObject( - topicId, MetadataObject.Type.TOPIC.name())), - () -> + topicId, MetadataObject.Type.TOPIC.name())); SessionUtils.doWithoutCommit( TagMetadataObjectRelMapper.class, mapper -> mapper.softDeleteTagMetadataObjectRelsByMetadataObject( - topicId, MetadataObject.Type.TOPIC.name())), - () -> + topicId, MetadataObject.Type.TOPIC.name())); SessionUtils.doWithoutCommit( StatisticMetaMapper.class, - mapper -> mapper.softDeleteStatisticsByEntityId(topicId)), - () -> + mapper -> mapper.softDeleteStatisticsByEntityId(topicId)); SessionUtils.doWithoutCommit( PolicyMetadataObjectRelMapper.class, mapper -> mapper.softDeletePolicyMetadataObjectRelsByMetadataObject( - topicId, MetadataObject.Type.TOPIC.name())), + topicId, MetadataObject.Type.TOPIC.name())); + } + }, () -> { - SessionUtils.doWithoutCommit( - EntityChangeLogMapper.class, - mapper -> - mapper.insertEntityChange( - metalakeName, - Entity.EntityType.TOPIC.name(), - topicFullName, - OperateType.DROP)); + if (deleteResult.get() > 0) { + SessionUtils.doWithoutCommit( + EntityChangeLogMapper.class, + mapper -> + mapper.insertEntityChange( + metalakeName, + Entity.EntityType.TOPIC.name(), + topicFullName, + OperateType.DROP)); + } }); - return true; + return deleteResult.get() > 0; } @Monitored( diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/service/ViewMetaService.java b/core/src/main/java/org/apache/gravitino/storage/relational/service/ViewMetaService.java index c65793abd2..177118c811 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/service/ViewMetaService.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/service/ViewMetaService.java @@ -258,6 +258,10 @@ public class ViewMetaService { mapper -> mapper.softDeleteTagMetadataObjectRelsByMetadataObject( viewId, MetadataObject.Type.VIEW.name())); + } + }, + () -> { + if (deleteResult.get() > 0) { SessionUtils.doWithoutCommit( EntityChangeLogMapper.class, mapper -> diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/session/SqlSessionFactoryHelper.java b/core/src/main/java/org/apache/gravitino/storage/relational/session/SqlSessionFactoryHelper.java index f96b7c7757..529dfad676 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/session/SqlSessionFactoryHelper.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/session/SqlSessionFactoryHelper.java @@ -31,7 +31,7 @@ import org.apache.gravitino.GravitinoEnv; import org.apache.gravitino.metrics.MetricsSystem; import org.apache.gravitino.metrics.source.RelationDatasourceMetricsSource; import org.apache.gravitino.storage.relational.JDBCBackend.JDBCBackendType; -import org.apache.gravitino.storage.relational.mapper.OperateTypeTypeHandler; +import org.apache.gravitino.storage.relational.mapper.OperateTypeHandler; import org.apache.gravitino.storage.relational.mapper.provider.MapperPackageProvider; import org.apache.gravitino.storage.relational.po.cache.OperateType; import org.apache.gravitino.utils.JdbcUrlUtils; @@ -112,9 +112,7 @@ public class SqlSessionFactoryHelper { // Initialize the configuration Configuration configuration = new Configuration(environment); configuration.setDatabaseId(jdbcType.name().toLowerCase()); - configuration - .getTypeHandlerRegistry() - .register(OperateType.class, new OperateTypeTypeHandler()); + configuration.getTypeHandlerRegistry().register(OperateType.class, new OperateTypeHandler()); ServiceLoader<MapperPackageProvider> loader = ServiceLoader.load(MapperPackageProvider.class); for (MapperPackageProvider provider : loader) { provider.getMapperClasses().forEach(configuration::addMapper); diff --git a/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestEntityChangeLogMapper.java b/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestEntityChangeLogMapper.java index 55f80f1790..8b4a3be8c2 100644 --- a/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestEntityChangeLogMapper.java +++ b/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestEntityChangeLogMapper.java @@ -22,6 +22,7 @@ package org.apache.gravitino.storage.relational.mapper.provider.base; import java.io.IOException; import java.nio.file.Files; import java.nio.file.Path; +import java.sql.PreparedStatement; import java.sql.SQLException; import java.sql.Statement; import java.util.List; @@ -123,7 +124,7 @@ public class TestEntityChangeLogMapper { @Test void testEntityChangeLogPruneOldEntries() throws SQLException { entityChangeLogMapper.insertEntityChange( - "metalake1", "SCHEMA", "metalake1.cat.schema", OperateType.INSERT); + "metalake1", "SCHEMA", "metalake1.cat.schema", OperateType.ALTER); forceCreatedAt("metalake1.cat.schema", 1000L); entityChangeLogMapper.insertEntityChange( "metalake1", "TABLE", "metalake1.cat.schema.tbl", OperateType.DROP); @@ -143,9 +144,9 @@ public class TestEntityChangeLogMapper { @Test void testEntityChangeLogSameTimestampOrderedById() throws SQLException { - entityChangeLogMapper.insertEntityChange("metalake1", "TABLE", "a", OperateType.INSERT); - entityChangeLogMapper.insertEntityChange("metalake1", "TABLE", "b", OperateType.INSERT); - entityChangeLogMapper.insertEntityChange("metalake1", "TABLE", "c", OperateType.INSERT); + entityChangeLogMapper.insertEntityChange("metalake1", "TABLE", "a", OperateType.ALTER); + entityChangeLogMapper.insertEntityChange("metalake1", "TABLE", "b", OperateType.ALTER); + entityChangeLogMapper.insertEntityChange("metalake1", "TABLE", "c", OperateType.ALTER); forceCreatedAt("a", 5_000_000L); forceCreatedAt("b", 5_000_000L); forceCreatedAt("c", 5_000_000L); @@ -157,13 +158,14 @@ public class TestEntityChangeLogMapper { } private void forceCreatedAt(String fullName, long createdAt) throws SQLException { - try (Statement statement = sharedSession.getConnection().createStatement()) { - statement.execute( - "UPDATE entity_change_log SET created_at = " - + createdAt - + " WHERE entity_full_name = '" - + fullName - + "'"); + try (PreparedStatement statement = + sharedSession + .getConnection() + .prepareStatement( + "UPDATE entity_change_log SET created_at = ? WHERE entity_full_name = ?")) { + statement.setLong(1, createdAt); + statement.setString(2, fullName); + statement.executeUpdate(); } sharedSession.clearCache(); }
