!1029 Memory Connector: Workers use Hetu Metastore to update tables based on a schedule

Merge pull request !1029 from arjunkrishna/worker-table-refresh-schedule
This commit is contained in:
i-robot 2021-08-06 01:31:23 +00:00 committed by Gitee
commit e9cee8ec6f
5 changed files with 37 additions and 3 deletions

View File

@ -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
--------

View File

@ -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获取最新的表的列表并在本地清理已经删除的表
## 示例
使用内存连接器创建表:

View File

@ -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)
{

View File

@ -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);

View File

@ -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<ConnectorViewDefinition> VIEW_CODEC = jsonCodec(ConnectorViewDefinition.class);
private static final JsonCodec<MemoryWriteTableHandle> 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<Long, TableInfo> tables = new ConcurrentHashMap<>();
private final Map<SchemaTableName, ConnectorViewDefinition> 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<CatalogEntity> 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<String> listSchemaNames(ConnectorSession session)
{