From ad0a5d0841ba9ab3806fab3917c161f8d146aa1b Mon Sep 17 00:00:00 2001 From: arjunkrishna Date: Thu, 5 Aug 2021 10:50:35 -0400 Subject: [PATCH] add scheduled refresh of tables --- hetu-docs/en/connector/memory.md | 5 ++++- hetu-docs/zh/connector/memory.md | 4 +++- .../plugin/memory/MemoryConnector.java | 5 +++++ .../plugin/memory/MemoryConnectorFactory.java | 5 ++++- .../plugin/memory/MemoryMetadata.java | 21 +++++++++++++++++++ 5 files changed, 37 insertions(+), 3 deletions(-) diff --git a/hetu-docs/en/connector/memory.md b/hetu-docs/en/connector/memory.md index cf46fb1c0..e830ae52a 100644 --- a/hetu-docs/en/connector/memory.md +++ b/hetu-docs/en/connector/memory.md @@ -28,7 +28,10 @@ memory.spill-path=/opt/hetu/data/spill distribute the data uniformly between all the workers. - Hetu Metastore must be configured. The default settings are included in `etc/hetu-metastore.properties`. - Check [Hetu Metastore](../admin/meta-store.md) for more information. + - Check [Hetu Metastore](../admin/meta-store.md) for more information. +- State Store must be configured to enable automatic memory refresh on workers. + - Check [State Store](../admin/state-store.md) for more information + - Automatic memory refresh will allow Memory Connector to clean unused tables more often resulting in more efficient use of memory. Examples -------- diff --git a/hetu-docs/zh/connector/memory.md b/hetu-docs/zh/connector/memory.md index f1377f355..0d7a01d22 100644 --- a/hetu-docs/zh/connector/memory.md +++ b/hetu-docs/zh/connector/memory.md @@ -21,7 +21,9 @@ memory.spill-path=/opt/hetu/data/spill - 关于更多详细信息与其他可选配置项,请参见**配置属性** 章节。 - 在`etc/config.properties`中,请确保`task.writer-count`的数字不小于配置的openLooKeng集群的节点个数。这会帮助把所有数据更均匀地分配到各个节点上。 - Hetu Metastore必须被妥善配置来保证内存连接器的正常功能。请参阅[Hetu Metastore](../admin/meta-store.md)。 - +- 必须配置StateStore使得worker上的表自动刷新生效 + - 请参阅[State Store](../admin/state-store.md) + - 自动刷新特性使得worker节点定期从metastore获取最新的表的列表,并在本地清理已经删除的表 ## 示例 使用内存连接器创建表: diff --git a/presto-memory/src/main/java/io/prestosql/plugin/memory/MemoryConnector.java b/presto-memory/src/main/java/io/prestosql/plugin/memory/MemoryConnector.java index 7025322aa..fb091c4c9 100644 --- a/presto-memory/src/main/java/io/prestosql/plugin/memory/MemoryConnector.java +++ b/presto-memory/src/main/java/io/prestosql/plugin/memory/MemoryConnector.java @@ -53,6 +53,11 @@ public class MemoryConnector this.tableProperties = ImmutableList.copyOf(requireNonNull(tableProperties, "tableProperties is null").getTableProperties()); } + void scheduleRefreshJob() + { + metadata.scheduleRefreshJob(); + } + @Override public ConnectorTransactionHandle beginTransaction(IsolationLevel isolationLevel, boolean readOnly) { diff --git a/presto-memory/src/main/java/io/prestosql/plugin/memory/MemoryConnectorFactory.java b/presto-memory/src/main/java/io/prestosql/plugin/memory/MemoryConnectorFactory.java index 4f1961502..31c1be64d 100644 --- a/presto-memory/src/main/java/io/prestosql/plugin/memory/MemoryConnectorFactory.java +++ b/presto-memory/src/main/java/io/prestosql/plugin/memory/MemoryConnectorFactory.java @@ -83,7 +83,10 @@ public class MemoryConnectorFactory MemoryThreadManager.initSharedThreadPool(injector.getInstance(MemoryConfig.class).getThreadPoolSize()); } - return injector.getInstance(MemoryConnector.class); + MemoryConnector memoryConnector = injector.getInstance(MemoryConnector.class); + memoryConnector.scheduleRefreshJob(); + + return memoryConnector; } catch (Exception e) { throwIfUnchecked(e); diff --git a/presto-memory/src/main/java/io/prestosql/plugin/memory/MemoryMetadata.java b/presto-memory/src/main/java/io/prestosql/plugin/memory/MemoryMetadata.java index 771555ed7..75f8e4530 100644 --- a/presto-memory/src/main/java/io/prestosql/plugin/memory/MemoryMetadata.java +++ b/presto-memory/src/main/java/io/prestosql/plugin/memory/MemoryMetadata.java @@ -63,6 +63,7 @@ import java.util.Map; import java.util.Optional; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; import java.util.stream.Collectors; import java.util.stream.Stream; @@ -88,16 +89,21 @@ public class MemoryMetadata public static final String NEXT_ID_KEY = "_NEXT_ID_"; // used as param key in CatalogEntity private static final JsonCodec VIEW_CODEC = jsonCodec(ConnectorViewDefinition.class); private static final JsonCodec OUTPUT_TABLE_HANDLE_JSON_CODEC = jsonCodec(MemoryWriteTableHandle.class); + private static final long metastoreCheckInterval = 10L; // time between scheduled metastore check to remove spilled data private final NodeManager nodeManager; private final TypeManager typeManager; private final AtomicLong nextTableId; private final HetuMetastore metastore; private final MemoryConfig config; + private final MemoryTableManager tableManager; + // tables and views are cached here private final Map tables = new ConcurrentHashMap<>(); private final Map views = new ConcurrentHashMap<>(); + private boolean refreshScheduled; // used to make sure that refresh is only scheduled once + @Inject public MemoryMetadata(TypeManager typeManager, NodeManager nodeManager, HetuMetastore metastore, MemoryTableManager tableManager, MemoryConfig memoryConfig) { @@ -105,6 +111,8 @@ public class MemoryMetadata this.nodeManager = requireNonNull(nodeManager, "nodeManager is null"); this.config = requireNonNull(memoryConfig, "memoryConfig is null"); this.metastore = metastore; + this.tableManager = tableManager; + this.refreshScheduled = false; Optional oldCatalog = metastore.getCatalog(MEM_KEY); if (!oldCatalog.isPresent()) { CatalogEntity newCatalog = CatalogEntity.builder() @@ -123,6 +131,19 @@ public class MemoryMetadata metastore.createDatabaseIfNotExist(databaseBuilder.build()); } + synchronized void scheduleRefreshJob() + { + if (!this.refreshScheduled && MemoryThreadManager.isSharedThreadPoolInitilized()) { + MemoryThreadManager.getSharedThreadPool().scheduleAtFixedRate(() -> { + if (tableManager != null) { + nextTableId.set(Long.parseLong(getMemoryCatalogEntity().getParameters().getOrDefault(NEXT_ID_KEY, nextTableId.toString()))); + tableManager.refreshTables(getTableIdSet(nextTableId.get())); + } + }, 1L, metastoreCheckInterval, TimeUnit.SECONDS); + this.refreshScheduled = true; + } + } + @Override public synchronized List listSchemaNames(ConnectorSession session) {