apache / apache/fluss

[Fluss/trino] Fluss Trino Engine Support

Open
#1,810 3 comments 3 reactions 0 assignees View on GitHub
Dominant language
Java
Stars
2.1k
Forks
625
Avg merge
3d 14h
Merged PRs (30d)
97

Description

### Search before asking

- [x] I searched in the [issues](https://github.com/apache/fluss/issues) and found nothing similar.

### Motivation

Technical Architecture Diagram

```mermaid
graph TB
subgraph "Trino Engine"
TrinoCoordinator["Trino Coordinator"]
TrinoWorker["Trino Workers"]
end

subgraph "Fluss-Trino Connector"
TrinoConnector["FlussTrinoConnector"]
TrinoMetadata["FlussTrinoMetadata"]
TrinoSplitManager["FlussTrinoSplitManager"]
TrinoPageSource["FlussTrinoPageSource"]
end

subgraph "Fluss Cluster"
CoordinatorServer["CoordinatorServer"]
TabletServers["TabletServers"]
end

subgraph "Lakehouse Storage"
PaimonTables["Paimon Tables"]
ObjectStorage["Object Storage"]
end

TrinoCoordinator --> TrinoConnector
TrinoWorker --> TrinoPageSource
TrinoMetadata --> CoordinatorServer
TrinoSplitManager --> TabletServers
TrinoPageSource --> TabletServers
TrinoPageSource --> PaimonTables
```

```mermaid
graph TB
subgraph "Trino 查询引擎层"
TC["Trino Coordinator
- 查询规划
- 任务调度
- 元数据管理"]
TW1["Trino Worker 1
- 查询执行
- 数据处理"]
TW2["Trino Worker 2
- 查询执行
- 数据处理"]
TWN["Trino Worker N
- 查询执行
- 数据处理"]
end

subgraph "Fluss-Trino 连接器层"
direction TB
FCF["FlussConnectorFactory
- 连接器工厂
- 配置管理"]
FC["FlussConnector
- 连接器实例
- 生命周期管理"]

subgraph "元数据管理"
FM["FlussMetadata
- 表元数据
- Schema 映射
- 分区信息"]
FTH["FlussTableHandle
- 表句柄
- 表标识符"]
FCH["FlussColumnHandle
- 列句柄
- 列元数据"]
end

subgraph "数据访问"
FSM["FlussSplitManager
- 数据分片管理
- 并行度控制"]
FPS["FlussPageSource
- 数据页面源
- 数据读取"]
FRS["FlussRecordSet
- 记录集
- 数据转换"]
end

subgraph "查询优化"
FPP["谓词下推
- 过滤条件
- 分区裁剪"]
FCP["列裁剪
- 投影下推
- I/O 优化"]
FAP["聚合下推
- COUNT/SUM
- 存储层计算"]
end
end

subgraph "Fluss 存储集群"
direction TB
CS["CoordinatorServer
- 集群协调
- 元数据管理
- 负载均衡"]

subgraph "TabletServer 集群"
TS1["TabletServer 1
- LogStore
- KvStore
- 数据分片"]
TS2["TabletServer 2
- LogStore
- KvStore
- 数据分片"]
TSN["TabletServer N
- LogStore
- KvStore
- 数据分片"]
end

subgraph "存储引擎"
LS["LogStore
- 流式数据
- Arrow 格式
- 列式存储"]
KS["KvStore
- 主键数据
- RocksDB
- 点查询"]
end
end

subgraph "分层存储"
direction TB
subgraph "热存储层"
LocalData["本地存储
- 最新数据
- 毫秒级延迟"]
end

subgraph "温存储层"
RemoteStorage["远程存储
- S3/HDFS/OSS
- 历史数据
- 成本优化"]
end

subgraph "分析存储层"
Lakehouse["Lakehouse 存储
- Apache Paimon
- Parquet 格式
- 分析优化"]
end
end

subgraph "外部系统"
ZK["ZooKeeper
- 集群协调
- 元数据存储"]
Catalog["数据目录
- Hive Metastore
- 元数据同步"]
end

%% 连接关系
TC --> FCF
TC --> FM
TW1 --> FPS
TW2 --> FPS
TWN --> FPS

FCF --> FC
FC --> FM
FC --> FSM
FC --> FPS

FM --> CS
FSM --> TS1
FSM --> TS2
FSM --> TSN
FPS --> LS
FPS --> KS

%% Union Read 支持
FPS -.->|"Union Read"| Lakehouse
FPS --> LocalData

%% 查询优化
FM --> FPP
FSM --> FCP
FPS --> FAP

%% 存储分层
LS --> RemoteStorage
KS --> RemoteStorage
LS -.->|"数据分层"| Lakehouse

%% 外部依赖
CS --> ZK
TS1 --> ZK
TS2 --> ZK
TSN --> ZK

CS --> Catalog
Lakehouse --> Catalog
```

Data Flow Diagram

```mermaid
sequenceDiagram
participant TC as "Trino Coordinator"
participant TW as "Trino Worker"
participant FC as "Fluss Connector"
participant CS as "CoordinatorServer"
participant TS as "TabletServer"
participant LS as "Lakehouse Storage"

Note over TC,LS: 查询执行流程

TC->>FC: 1. 获取表元数据
FC->>CS: 2. 请求表信息
CS-->>FC: 3. 返回表元数据
FC-->>TC: 4. 返回表句柄

TC->>FC: 5. 生成执行计划
FC->>FC: 6. 应用谓词下推
FC->>FC: 7. 应用列裁剪

TC->>TW: 8. 分发查询任务
TW->>FC: 9. 请求数据分片
FC->>TS: 10. 读取实时数据
FC->>LS: 11. 读取历史数据 (Union Read)

TS-->>FC: 12. 返回实时数据
LS-->>FC: 13. 返回历史数据
FC->>FC: 14. 数据合并与转换
FC-->>TW: 15. 返回结果页面
TW-->>TC: 16. 返回查询结果
```

Module Dependency Graph

```mermaid
graph LR
subgraph "Maven 模块结构"
FT["fluss-trino
(父模块)"]
FTC["fluss-trino-common
(通用组件)"]
FT435["fluss-trino-435
(版本特定)"]
FT436["fluss-trino-436
(版本特定)"]
FT437["fluss-trino-437
(版本特定)"]
end

subgraph "Fluss 核心模块"
FC["fluss-client"]
FCom["fluss-common"]
FS["fluss-server"]
end

subgraph "Trino SPI"
TSPI["trino-spi"]
TTest["trino-testing"]
end

FT --> FTC
FT --> FT435
FT --> FT436
FT --> FT437

FTC --> FC
FTC --> FCom
FTC --> TSPI

FT435 --> FTC
FT435 --> TSPI
FT436 --> FTC
FT436 --> TSPI
FT437 --> FTC
FT437 --> TSPI

FTC -.->|"测试依赖"| FS
FTC -.->|"测试依赖"| TTest
```

### Overall Architecture Design

Fluss currently only supports Flink as the OLAP engine, but Trino support has been planned in the roadmap. Based on the existing distributed architecture, an independent connector module needs to be developed for Trino.

### Implementation Plan

Phase 1: Create Trino Connector Module
1. Module Structure: Refer to Trino Multi-Version Support Mode
2. Connector Factory: FlussConnectorFactory serves as the Trino plugin entry point
3. Dependency Injection: Using Guice to Manage Components
Referring to the structure of the existing Flink integration module, create a new module:
```
fluss-trino/
├── fluss-trino-common/ # 通用组件
├── fluss-trino-{version}/ # 版本特定实现
└── pom.xml
```

1. Root Module POM Configuration
```


4.0.0

com.alibaba.fluss
fluss
0.8-SNAPSHOT


fluss-trino
pom

Fluss : Engine Trino


fluss-trino-common
fluss-trino-435
fluss-trino-436
fluss-trino-437



435

```
2. General Module POM Configuration
fluss-trino/fluss-trino-common/pom.xml
```


4.0.0

com.alibaba.fluss
fluss-trino
0.8-SNAPSHOT


fluss-trino-common

Fluss : Engine Trino : Common




com.alibaba.fluss
fluss-client
${project.version}



com.alibaba.fluss
fluss-common
${project.version}




io.trino
trino-spi
${trino.version}
provided




com.fasterxml.jackson.core
jackson-databind
provided




com.google.guava
guava
provided




com.alibaba.fluss
fluss-test-utils
test



io.trino
trino-testing
${trino.version}
test


```
3. Version-Specific Module POM Configuration
fluss-trino/fluss-trino-435/pom.xml
```


4.0.0

com.alibaba.fluss
fluss-trino
0.8-SNAPSHOT


fluss-trino-435

Fluss : Engine Trino : 435


435
435





com.alibaba.fluss
fluss-trino-common
${project.version}


*
*






io.trino
trino-spi
${trino.minor.version}
provided




com.alibaba.fluss
fluss-trino-common
${project.version}
test
test-jar



com.alibaba.fluss
fluss-server
${project.version}
test



io.trino
trino-testing
${trino.minor.version}
test







org.apache.maven.plugins
maven-shade-plugin


shade-fluss
package

shade


false


com.alibaba.fluss:*





com.alibaba.fluss
com.alibaba.fluss.shaded.fluss








```
### Guice Dependency Injection Module
fluss-trino/fluss-trino-common/src/main/java/com/alibaba/fluss/trino/FlussConnectorModule.java
```
package com.alibaba.fluss.trino;

import com.alibaba.fluss.config.Configuration;
import com.google.inject.Binder;
import com.google.inject.Module;
import com.google.inject.Scopes;
import io.trino.spi.type.TypeManager;

import static io.airlift.configuration.ConfigBinder.configBinder;
import static java.util.Objects.requireNonNull;

/**
* Guice module for Fluss connector dependency injection.
*/
public class FlussConnectorModule implements Module {

private final String catalogName;
private final Configuration flussConfig;
private final TypeManager typeManager;

public FlussConnectorModule(String catalogName, Configuration flussConfig, TypeManager typeManager) {
this.catalogName = requireNonNull(catalogName, "catalogName is null");
this.flussConfig = requireNonNull(flussConfig, "flussConfig is null");
this.typeManager = requireNonNull(typeManager, "typeManager is null");
}

@Override
public void configure(Binder binder) {
// Bind configuration
binder.bind(String.class).annotatedWith(ForFlussConnector.class).toInstance(catalogName);
binder.bind(Configuration.class).toInstance(flussConfig);
binder.bind(TypeManager.class).toInstance(typeManager);

// Bind connector components
binder.bind(FlussConnector.class).in(Scopes.SINGLETON);
binder.bind(FlussMetadata.class).in(Scopes.SINGLETON);
binder.bind(FlussSplitManager.class).in(Scopes.SINGLETON);
binder.bind(FlussPageSourceProvider.class).in(Scopes.SINGLETON);
binder.bind(FlussTransactionManager.class).in(Scopes.SINGLETON);

// Bind client factory
binder.bind(FlussClientManager.class).in(Scopes.SINGLETON);

// Configure connector config
configBinder(binder).bindConfig(FlussConnectorConfig.class);
}
}
```
### Table Handle
fluss-trino/fluss-trino-common/src/main/java/com/alibaba/fluss/trino/FlussTableHandle.java
```
package com.alibaba.fluss.trino;

import com.alibaba.fluss.metadata.TableInfo;
import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonProperty;
import io.trino.spi.connector.ConnectorTableHandle;

import static java.util.Objects.requireNonNull;

/**
* Table handle for Fluss connector.
*/
public class FlussTableHandle implements ConnectorTableHandle {

private final String schemaName;
private final String tableName;
private final TableInfo tableInfo;

@JsonCreator
public FlussTableHandle(
@JsonProperty("schemaName") String schemaName,
@JsonProperty("tableName") String tableName,
@JsonProperty("tableInfo") TableInfo tableInfo) {
this.schemaName = requireNonNull(schemaName, "schemaName is null");
this.tableName = requireNonNull(tableName, "tableName is null");
this.tableInfo = requireNonNull(tableInfo, "tableInfo is null");
}

@JsonProperty
public String getSchemaName() {
return schemaName;
}

@JsonProperty
public String getTableName() {
return tableName;
}

@JsonProperty
public TableInfo getTableInfo() {
return tableInfo;
}

@Override
public boolean equals(Object obj) {
if (this == obj) {
return true;
}
if (obj == null || getClass() != obj.getClass()) {
return false;
}
FlussTableHandle that = (FlussTableHandle) obj;
return schemaName.equals(that.schemaName) &&
tableName.equals(that.tableName);
}

@Override
public int hashCode() {
return Objects.hash(schemaName, tableName);
}

@Override
public String toString() {
return "FlussTableHandle{" +
"schemaName='" + schemaName + '\'' +
", tableName='" + tableName + '\'' +
", tableInfo=" + tableInfo +
'}';
}
}
```

Column Handle
fluss-trino/fluss-trino-common/src/main/java/com/alibaba/fluss/trino/FlussColumnHandle.java

```
package com.alibaba.fluss.trino;

import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonProperty;
import io.trino.spi.connector.ColumnHandle;
import io.trino.spi.type.Type;

import java.util.Objects;

import static java.util.Objects.requireNonNull;

/**
* Column handle for Fluss connector.
*/
public class FlussColumnHandle implements ColumnHandle {

private final String name;
private final Type type;
private final int ordinalPosition;

@JsonCreator
public FlussColumnHandle(
@JsonProperty("name") String name,
@JsonProperty("type") Type type,
@JsonProperty("ordinalPosition") int ordinalPosition) {
this.name = requireNonNull(name, "name is null");
this.type = requireNonNull(type, "type is null");
this.ordinalPosition = ordinalPosition;
}

@JsonProperty
public String getName() {
return name;
}

@JsonProperty
public Type getType() {
return type;
}

@JsonProperty
public int getOrdinalPosition() {
return ordinalPosition;
}

@Override
public boolean equals(Object obj) {
if (this == obj) {
return true;
}
if (obj == null || getClass() != obj.getClass()) {
return false;
}
FlussColumnHandle that = (FlussColumnHandle) obj;
return ordinalPosition == that.ordinalPosition &&
Objects.equals(name, that.name) &&
Objects.equals(type, that.type);
}

@Override
public int hashCode() {
return Objects.hash(name, type, ordinalPosition);
}

@Override
public String toString() {
return "FlussColumnHandle{" +
"name='" + name + '\'' +
", type=" + type +
", ordinalPosition=" + ordinalPosition +
'}';
}
}
```
### Metadata Manager
fluss-trino/fluss-trino-common/src/main/java/com/alibaba/fluss/trino/FlussMetadata.java

```
package com.alibaba.fluss.trino;

import com.alibaba.fluss.client.admin.Admin;
import com.alibaba.fluss.metadata.TableInfo;
import com.alibaba.fluss.metadata.TablePath;
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableMap;
import io.airlift.log.Logger;
import io.trino.spi.connector.ColumnHandle;
import io.trino.spi.connector.ColumnMetadata;
import io.trino.spi.connector.ConnectorMetadata;
import io.trino.spi.connector.ConnectorSession;
import io.trino.spi.connector.ConnectorTableHandle;
import io.trino.spi.connector.ConnectorTableMetadata;
import io.trino.spi.connector.SchemaTableName;
import io.trino.spi.connector.SchemaTablePrefix;
import io.trino.spi.type.TypeManager;

import javax.inject.Inject;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.ExecutionException;

import static java.util.Objects.requireNonNull;

/**
* Metadata provider for Fluss connector.
*/
public class FlussMetadata implements ConnectorMetadata {

private static final Logger log = Logger.get(FlussMetadata.class);

private final FlussClientManager clientManager;
private final TypeManager typeManager;
private final FlussConnectorConfig config;

@Inject
public FlussMetadata(
FlussClientManager clientManager,
TypeManager typeManager,
FlussConnectorConfig config) {
this.clientManager = requireNonNull(clientManager, "clientManager is null");
this.typeManager = requireNonNull(typeManager, "typeManager is null");
this.config = requireNonNull(config, "config is null");
}

@Override
public List listSchemaNames(ConnectorSession session) {
try {
Admin admin = clientManager.getAdmin();
return admin.listDatabases().get();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException("Interrupted while listing databases", e);
} catch (ExecutionException e) {
throw new RuntimeException("Failed to list databases", e);
}
}

@Override
public List listTables(ConnectorSession session, Optional schemaName) {
try {
Admin admin = clientManager.getAdmin();
ImmutableList.Builder tables = ImmutableList.builder();

List databases = schemaName.isPresent()
? ImmutableList.of(schemaName.get())
: admin.listDatabases().get();

for (String database : databases) {
List tableNames = admin.listTables(database).get();
for (String tableName : tableNames) {
tables.add(new SchemaTableName(database, tableName));
}
}

return tables.build();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException("Interrupted while listing tables", e);
} catch (ExecutionException e) {
throw new RuntimeException("Failed to list tables", e);
}
}

@Override
public ConnectorTableHandle getTableHandle(ConnectorSession session, SchemaTableName tableName) {
try {
Admin admin = clientManager.getAdmin();
TablePath tablePath = new TablePath(tableName.getSchemaName(), tableName.getTableName());

TableInfo tableInfo = admin.getTableInfo(tablePath).get();
if (tableInfo == null) {
return null;
}

return new FlussTableHandle(
tableName.getSchemaName(),
tableName.getTableName(),
tableInfo);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException("Interrupted while getting table handle", e);
} catch (ExecutionException e) {
log.debug(e, "Table %s not found", tableName);
return null;
}
}

@Override
public ConnectorTableMetadata getTableMetadata(ConnectorSession session, ConnectorTableHandle table) {
FlussTableHandle flussTable = (FlussTableHandle) table;

ImmutableList.Builder columns = ImmutableList.builder();
List fieldNames = flussTable.getTableInfo().getRowType().getFieldNames();

for (int i = 0; i < fieldNames.size(); i++) {
String fieldName = fieldNames.get(i);
com.alibaba.fluss.types.DataType flussType = flussTable.getTableInfo().getRowType().getTypeAt(i);
io.trino.spi.type.Type trinoType = FlussTypeUtils.toTrinoType(flussType, typeManager);

columns.add(ColumnMetadata.builder()
.setName(fieldName)
.setType(trinoType)
.build());
}

return new ConnectorTableMetadata(
new SchemaTableName(flussTable.getSchemaName(), flussTable.getTableName()),
columns.build());
}

@Override
public Map getColumnHandles(ConnectorSession session, ConnectorTableHandle tableHandle) {
FlussTableHandle flussTable = (FlussTableHandle) tableHandle;

ImmutableMap.Builder columnHandles = ImmutableMap.builder();
List fieldNames = flussTable.getTableInfo().getRowType().getFieldNames();

for (int i = 0; i < fieldNames.size(); i++) {
String fieldName = fieldNames.get(i);
com.alibaba.fluss.types.DataType flussType = flussTable.getTableInfo().getRowType().getTypeAt(i);
io.trino.spi.type.Type trinoType = FlussTypeUtils.toTrinoType(flussType, typeManager);

columnHandles.put(fieldName, new FlussColumnHandle(fieldName, trinoType, i));
}

return columnHandles.build();
}

@Override
public ColumnMetadata getColumnMetadata(ConnectorSession session, ConnectorTableHandle tableHandle, ColumnHandle columnHandle) {
FlussColumnHandle flussColumn = (FlussColumnHandle) columnHandle;
return ColumnMetadata.builder()
.setName(flussColumn.getName())
.setType(flussColumn.getType())
.build();
}
}
```
### Transaction Manager
fluss-trino/fluss-trino-common/src/main/java/com/alibaba/fluss/trino/FlussTransactionManager.java

```
package com.alibaba.fluss.trino;

import io.trino.spi.connector.ConnectorTransactionHandle;
import io.trino.spi.transaction.IsolationLevel;

import javax.inject.Inject;

/**
* Transaction manager for Fluss connector.
* Currently supports read-only transactions.
*/
public class FlussTransactionManager {

@Inject
public FlussTransactionManager() {
}

public ConnectorTransactionHandle beginTransaction(IsolationLevel isolationLevel, boolean readOnly, boolean autoCommit) {
// For now, we only support read-only transactions
if (!readOnly) {
throw new UnsupportedOperationException("Write transactions are not yet supported");
}

return new FlussTransactionHandle();
}

public void commit(ConnectorTransactionHandle transaction) {
// No-op for read-only transactions
}

public void rollback(ConnectorTransactionHandle transaction) {
// No-op for read-only transactions
}
}
```

### Transaction Handle
fluss-trino/fluss-trino-common/src/main/java/com/alibaba/fluss/trino/FlussTransactionHandle.java

```
package com.alibaba.fluss.trino;

import com.fasterxml.jackson.annotation.JsonCreator;
import io.trino.spi.connector.ConnectorTransactionHandle;

import java.util.UUID;

/**
* Transaction handle for Fluss connector.
*/
public class FlussTransactionHandle implements ConnectorTransactionHandle {

private final String transactionId;

@JsonCreator
public FlussTransactionHandle() {
this.transactionId = UUID.randomUUID().toString();
}

public String getTransactionId() {
return transactionId;
}

@Override
public boolean equals(Object obj) {
if (this == obj) {
return true;
}
if (obj == null || getClass() != obj.getClass()) {
return false;
}
FlussTransactionHandle that = (FlussTransactionHandle) obj;
return transactionId.equals(that.transactionId);
}

@Override
public int hashCode() {
return transactionId.hashCode();
}

@Override
public String toString() {
return "FlussTransactionHandle{" +
"transactionId='" + transactionId + '\'' +
'}';
}
}
```
### Sharding Manager
fluss-trino/fluss-trino-common/src/main/java/com/alibaba/fluss/trino/FlussSplitManager.java
```
package com.alibaba.fluss.trino;

import com.alibaba.fluss.metadata.TableBucket;
import com.alibaba.fluss.metadata.TableInfo;
import com.alibaba.fluss.metadata.TablePath;
import com.google.common.collect.ImmutableList;
import io.trino.spi.connector.ConnectorSession;
import io.trino.spi.connector.ConnectorSplit;
import io.trino.spi.connector.ConnectorSplitManager;
import io.trino.spi.connector.ConnectorSplitSource;
import io.trino.spi.connector.ConnectorTableHandle;
import io.trino.spi.connector.ConnectorTransactionHandle;
import io.trino.spi.connector.Constraint;
import io.trino.spi.connector.DynamicFilter;
import io.trino.spi.connector.FixedSplitSource;

import javax.inject.Inject;
import java.util.List;

import static java.util.Objects.requireNonNull;

/**
* Split manager for Fluss connector.
*/
public class FlussSplitManager implements ConnectorSplitManager {

private final FlussClientManager clientManager;

@Inject
public FlussSplitManager(FlussClientManager clientManager) {
this.clientManager = requireNonNull(clientManager, "clientManager is null");
}

@Override
public ConnectorSplitSource getSplits(
ConnectorTransactionHandle transaction,
ConnectorSession session,
ConnectorTableHandle tableHandle,
DynamicFilter dynamicFilter,
Constraint constraint) {

FlussTableHandle flussTable = (FlussTableHandle) tableHandle;
TableInfo tableInfo = flussTable.getTableInfo();
TablePath tablePath = new TablePath(flussTable.getSchemaName(), flussTable.getTableName());

// Create splits based on table buckets
ImmutableList.Builder splits = ImmutableList.builder();

// Get bucket count from table info
int bucketCount = tableInfo.getTableDescriptor().getDistribution().getBucketCount().orElse(1);

for (int bucketId = 0; bucketId < bucketCount; bucketId++) {
TableBucket tableBucket = new TableBucket(tableInfo.getTableId(), bucketId);
splits.add(new FlussSplit(tablePath, tableBucket));
}

return new FixedSplitSource(splits.build());
}
}
```
### Sharding Class
fluss-trino/fluss-trino-common/src/main/java/com/alibaba/fluss/trino/FlussSplit.java

```
package com.alibaba.fluss.trino;

import com.alibaba.fluss.metadata.TableBucket;
import com.alibaba.fluss.metadata.TablePath;
import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonProperty;
import io.trino.spi.HostAddress;
import io.trino.spi.connector.ConnectorSplit;

import java.util.List;
import java.util.Objects;

import static java.util.Objects.requireNonNull;

/**
* Split for Fluss connector representing a table bucket.
*/
public class FlussSplit implements ConnectorSplit {

private final TablePath tablePath;
private final TableBucket tableBucket;

@JsonCreator
public FlussSplit(
@JsonProperty("tablePath") TablePath tablePath,
@JsonProperty("tableBucket") TableBucket tableBucket) {
this.tablePath = requireNonNull(tablePath, "tablePath is null");
this.tableBucket = requireNonNull(tableBucket, "tableBucket is null");
}

@JsonProperty
public TablePath getTablePath() {
return tablePath;
}

@JsonProperty
public TableBucket getTableBucket() {
return tableBucket;
}

@Override
public boolean isRemotelyAccessible() {
return true;
}

@Override
public List getAddresses() {
// Return empty list for now - Trino will handle scheduling
return List.of();
}

@Override
public Object getInfo() {
return this;
}

@Override
public boolean equals(Object obj) {
if (this == obj) {
return true;
}
if (obj == null || getClass() != obj.getClass()) {
return false;
}
FlussSplit that = (FlussSplit) obj;
return Objects.equals(tablePath, that.tablePath) &&
Objects.equals(tableBucket, that.tableBucket);
}

@Override
public int hashCode() {
return Objects.hash(tablePath, tableBucket);
}

@Override
public String toString() {
return "FlussSplit{" +
"tablePath=" + tablePath +
", tableBucket=" + tableBucket +
'}';
}
}
```

Phase 2: Core Connector Component Development

1. Trino Connector Factory
Based on the DynamicTableSourceFactory and DynamicTableSinkFactory patterns of Flink, implement a similar factory class for Trino:

- FlussTrinoConnectorFactory: Main factory class
- FlussTrinoTableHandle: Table metadata processing
- FlussTrinoColumnHandle: Column metadata processing

Implementation of Core Components in the Second Phase

1. FlussTrinoConnectorFactory

This is the enhanced version of the connector factory for the second phase, which is extended based on the FlussConnectorFactory from the first phase:
```
/*
* 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 com.alibaba.fluss.trino;

import com.alibaba.fluss.config.Configuration;
import com.alibaba.fluss.metadata.TablePath;
import com.google.inject.Injector;
import io.airlift.bootstrap.Bootstrap;
import io.airlift.log.Logger;
import io.trino.spi.connector.Connector;
import io.trino.spi.connector.ConnectorContext;
import io.trino.spi.connector.ConnectorFactory;

import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;

import static java.util.Objects.requireNonNull;

/**
* Enhanced Fluss Trino connector factory for second phase implementation.
* Supports advanced features like Union Read, predicate pushdown, and column pruning.
*/
public class FlussTrinoConnectorFactory implements ConnectorFactory {

private static final Logger log = Logger.get(FlussTrinoConnectorFactory.class);

// Cache for connector instances to support connection pooling
private final Map connectorCache = new ConcurrentHashMap<>();

@Override
public String getName() {
return "fluss";
}

@Override
public Connector create(String catalogName, Map config, ConnectorContext context) {
requireNonNull(catalogName, "catalogName is null");
requireNonNull(config, "config is null");
requireNonNull(context, "context is null");

// Check if connector already exists for this catalog
return connectorCache.computeIfAbsent(catalogName, name -> createNewConnector(name, config, context));
}

private Connector createNewConnector(String catalogName, Map config, ConnectorContext context) {
try {
log.info("Creating new Fluss Trino connector for catalog: %s", catalogName);

// Convert Trino config to Fluss configuration with enhanced validation
Configuration flussConfig = createFlussConfiguration(config);

// Validate connection to Fluss cluster
validateFlussConnection(flussConfig);

// Create enhanced Guice injector with all necessary modules
Bootstrap app = new Bootstrap(
new FlussTrinoConnectorModule(catalogName, flussConfig, context.getTypeManager()),
new FlussTrinoMetadataModule(),
new FlussTrinoOptimizationModule()
);

Injector injector = app
.doNotInitializeLogging()
.initialize();

Connector connector = injector.getInstance(FlussTrinoConnector.class);
log.info("Successfully created Fluss Trino connector for catalog: %s", catalogName);

return connector;
} catch (Exception e) {
log.error(e, "Failed to create Fluss Trino connector for catalog: %s", catalogName);
throw new RuntimeException("Failed to create Fluss Trino connector", e);
}
}

private Configuration createFlussConfiguration(Map config) {
Configuration flussConfig = new Configuration();

// Bootstrap servers (required)
String bootstrapServers = config.get("bootstrap.servers");
if (bootstrapServers == null || bootstrapServers.trim().isEmpty()) {
throw new IllegalArgumentException("bootstrap.servers is required and cannot be empty");
}
flussConfig.setString("bootstrap.servers", bootstrapServers);

// Enhanced configuration mapping with validation
config.forEach((key, value) -> {
if (key.startsWith("fluss.")) {
// Remove fluss. prefix for internal configuration
String flussKey = key.substring(6);
flussConfig.setString(flussKey, value);
} else if (key.startsWith("trino.fluss.")) {
// Trino-specific Fluss configurations
String flussKey = key.substring(12);
flussConfig.setString(flussKey, value);
}
});

// Set default configurations for optimal performance
setDefaultConfigurations(flussConfig, config);

return flussConfig;
}

private void setDefaultConfigurations(Configuration flussConfig, Map config) {
// Connection pool settings
flussConfig.setString("client.connection.max-idle-time",
config.getOrDefault("connection.max-idle-time", "10min"));
flussConfig.setString("client.request.timeout",
config.getOrDefault("request.timeout", "60s"));

// Performance optimizations
flussConfig.setString("client.scanner.fetch.max-wait-time",
config.getOrDefault("scanner.fetch.max-wait-time", "500ms"));
flussConfig.setString("client.scanner.fetch.min-bytes",
config.getOrDefault("scanner.fetch.min-bytes", "1MB"));

// Enable Union Read by default for better analytics performance
flussConfig.setString("client.union.read.enabled",
config.getOrDefault("union.read.enabled", "true"));

// Enable column pruning for better I/O performance
flussConfig.setString("client.column.pruning.enabled",
config.getOrDefault("column.pruning.enabled", "true"));
}

private void validateFlussConnection(Configuration flussConfig) {
try {
// Basic connection validation - this would be expanded in a real implementation
String bootstrapServers = flussConfig.getString("bootstrap.servers");
if (bootstrapServers == null) {
throw new IllegalArgumentException("Invalid Fluss configuration: bootstrap.servers is null");
}

// Additional validation logic would go here
log.info("Fluss connection configuration validated successfully");
} catch (Exception e) {
throw new RuntimeException("Failed to validate Fluss connection configuration", e);
}
}

/**
* Shutdown method to clean up resources
*/
public void shutdown() {
log.info("Shutting down Fluss Trino connector factory");
connectorCache.values().forEach(connector -> {
try {
connector.shutdown();
} catch (Exception e) {
log.warn(e, "Error shutting down connector");
}
});
connectorCache.clear();
}
}
```
2. FlussTrinoTableHandle

Enhanced table handle, supporting predicate pushdown and partition information:

```
package com.alibaba.fluss.trino;

import com.alibaba.fluss.metadata.TableInfo;
import com.alibaba.fluss.metadata.TablePath;
import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableSet;
import io.trino.spi.connector.ColumnHandle;
import io.trino.spi.connector.ConnectorTableHandle;
import io.trino.spi.predicate.TupleDomain;

import java.util.List;
import java.util.Objects;
import java.util.Optional;
import java.util.Set;

import static java.util.Objects.requireNonNull;

/**
* Enhanced table handle for Fluss Trino connector with support for:
* - Predicate pushdown
* - Column pruning
* - Partition pruning
* - Union Read optimization
*/
public class FlussTrinoTableHandle implements ConnectorTableHandle {

private final String schemaName;
private final String tableName;
private final TableInfo tableInfo;
private final TablePath tablePath;

// Enhanced features for second phase
private final TupleDomain constraint;
private final Optional> projectedColumns;
private final Optional> partitionFilters;
private final boolean unionReadEnabled;
private final Optional limit;

@JsonCreator
public FlussTrinoTableHandle(
@JsonProperty("schemaName") String schemaName,
@JsonProperty("tableName") String tableName,
@JsonProperty("tableInfo") TableInfo tableInfo,
@JsonProperty("constraint") TupleDomain constraint,
@JsonProperty("projectedColumns") Optional> projectedColumns,
@JsonProperty("partitionFilters") Optional> partitionFilters,
@JsonProperty("unionReadEnabled") boolean unionReadEnabled,
@JsonProperty("limit") Optional limit) {
this.schemaName = requireNonNull(schemaName, "schemaName is null");
this.tableName = requireNonNull(tableName, "tableName is null");
this.tableInfo = requireNonNull(tableInfo, "tableInfo is null");
this.tablePath = new TablePath(schemaName, tableName);
this.constraint = requireNonNull(constraint, "constraint is null");
this.projectedColumns = requireNonNull(projectedColumns, "projectedColumns is null");
this.partitionFilters = requireNonNull(partitionFilters, "partitionFilters is null");
this.unionReadEnabled = unionReadEnabled;
this.limit = requireNonNull(limit, "limit is null");
}

// Convenience constructor for basic table handle
public FlussTrinoTableHandle(String schemaName, String tableName, TableInfo tableInfo) {
this(schemaName, tableName, tableInfo,
TupleDomain.all(),
Optional.empty(),
Optional.empty(),
true,
Optional.empty());
}

@JsonProperty
public String getSchemaName() {
return schemaName;
}

@JsonProperty
public String getTableName() {
return tableName;
}

@JsonProperty
public TableInfo getTableInfo() {
return tableInfo;
}

public TablePath getTablePath() {
return tablePath;
}

@JsonProperty
public TupleDomain getConstraint() {
return constraint;
}

@JsonProperty
public Optional> getProjectedColumns() {
return projectedColumns;
}

@JsonProperty
public Optional> getPartitionFilters() {
return partitionFilters;
}

@JsonProperty
public boolean isUnionReadEnabled() {
return unionReadEnabled;
}

@JsonProperty
public Optional getLimit() {
return limit;
}

/**
* Create a new table handle with applied constraint (predicate pushdown)
*/
public FlussTrinoTableHandle withConstraint(TupleDomain newConstraint) {
return new FlussTrinoTableHandle(
schemaName, tableName, tableInfo,
newConstraint, projectedColumns, partitionFilters,
unionReadEnabled, limit);
}

/**
* Create a new table handle with projected columns (column pruning)
*/
public FlussTrinoTableHandle withProjectedColumns(Set columns) {
return new FlussTrinoTableHandle(
schemaName, tableName, tableInfo,
constraint, Optional.of(ImmutableSet.copyOf(columns)), partitionFilters,
unionReadEnabled, limit);
}

/**
* Create a new table handle with partition filters (partition pruning)
*/
public FlussTrinoTableHandle withPartitionFilters(List filters) {
return new FlussTrinoTableHandle(
schemaName, tableName, tableInfo,
constraint, projectedColumns, Optional.of(ImmutableList.copyOf(filters)),
unionReadEnabled, limit);
}

/**
* Create a new table handle with limit pushdown
*/
public FlussTrinoTableHandle withLimit(long limitValue) {
return new FlussTrinoTableHandle(
schemaName, tableName, tableInfo,
constraint, projectedColumns, partitionFilters,
unionReadEnabled, Optional.of(limitValue));
}

/**
* Create a new table handle with Union Read setting
*/
public FlussTrinoTableHandle withUnionRead(boolean enabled) {
return new FlussTrinoTableHandle(
schemaName, tableName, tableInfo,
constraint, projectedColumns, partitionFilters,
enabled, limit);
}

/**
* Check if this table handle has any optimizations applied
*/
public boolean hasOptimizations() {
return !constraint.isAll() ||
projectedColumns.isPresent() ||
partitionFilters.isPresent() ||
limit.isPresent();
}

/**
* Get the effective columns to read (considering column pruning)
*/
public Optional> getEffectiveColumns() {
return projectedColumns;
}

@Override
public boolean equals(Object obj) {
if (this == obj) {
return true;
}
if (obj == null || getClass() != obj.getClass()) {
return false;
}
FlussTrinoTableHandle that = (FlussTrinoTableHandle) obj;
return Objects.equals(schemaName, that.schemaName) &&
Objects.equals(tableName, that.tableName) &&
Objects.equals(constraint, that.constraint) &&
Objects.equals(projectedColumns, that.projectedColumns) &&
Objects.equals(partitionFilters, that.partitionFilters) &&
unionReadEnabled == that.unionReadEnabled &&
Objects.equals(limit, that.limit);
}

@Override
public int hashCode() {
return Objects.hash(schemaName, tableName, constraint, projectedColumns,
partitionFilters, unionReadEnabled, limit);
}

@Override
public String toString() {
return "FlussTrinoTableHandle{" +
"schemaName='" + schemaName + '\'' +
", tableName='" + tableName + '\'' +
", constraint=" + constraint +
", projectedColumns=" + projectedColumns +
", partitionFilters=" + partitionFilters +
", unionReadEnabled=" + unionReadEnabled +
", limit=" + limit +
'}';
}
}
```

3.FlussTrinoColumnHandle
Enhanced column handle, supporting more metadata information:

```
package com.alibaba.fluss.trino;

import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonProperty;
import io.trino.spi.connector.ColumnHandle;
import io.trino.spi.type.Type;

import java.util.Objects;
import java.util.Optional;

import static java.util.Objects.requireNonNull;

/**
* Enhanced column handle for Fluss Trino connector with support for:
* - Column metadata
* - Partition key information
* - Primary key information
* - Column pruning optimization
*/
public class FlussTrinoColumnHandle implements ColumnHandle {

private final String name;
private final Type type;
private final int ordinalPosition;
private final boolean isPartitionKey;
private final boolean isPrimaryKey;
private final boolean isNullable;
private final Optional comment;
private final Optional defaultValue;

@JsonCreator
public FlussTrinoColumnHandle(
@JsonProperty("name") String name,
@JsonProperty("type") Type type,
@JsonProperty("ordinalPosition") int ordinalPosition,
@JsonProperty("isPartitionKey") boolean isPartitionKey,
@JsonProperty("isPrimaryKey") boolean isPrimaryKey,
@JsonProperty("isNullable") boolean isNullable,
@JsonProperty("comment") Optional comment,
@JsonProperty("defaultValue") Optional defaultValue) {
this.name = requireNonNull(name, "name is null");
this.type = requireNonNull(type, "type is null");
this.ordinalPosition = ordinalPosition;
this.isPartitionKey = isPartitionKey;
this.isPrimaryKey = isPrimaryKey;
this.isNullable = isNullable;
this.comment = requireNonNull(comment, "comment is null");
this.defaultValue = requireNonNull(defaultValue, "defaultValue is null");
}

// Convenience constructor for basic column handle
public FlussTrinoColumnHandle(String name, Type type, int ordinalPosition) {
this(name, type, ordinalPosition, false, false, true,
Optional.empty(), Optional.empty());
}

@JsonProperty
public String getName() {
return name;
}

@JsonProperty
public Type getType() {
return type;
}

@JsonProperty
public int getOrdinalPosition() {
return ordinalPosition;
}

@JsonProperty
public boolean isPartitionKey() {
return isPartitionKey;
}

@JsonProperty
public boolean isPrimaryKey() {
return isPrimaryKey;
}

@JsonProperty
public boolean isNullable() {
return isNullable;
}

@JsonProperty
public Optional getComment() {
return comment;
}

@JsonProperty
public Optional getDefaultValue() {
return defaultValue;
}

/**
* Create a new column handle with partition key flag
*/
public FlussTrinoColumnHandle withPartitionKey(boolean partitionKey) {
return new FlussTrinoColumnHandle(
name, type, ordinalPosition, partitionKey, isPrimaryKey,
isNullable, comment, defaultValue);
}

/**
* Create a new column handle with primary key flag
*/
public FlussTrinoColumnHandle withPrimaryKey(boolean primaryKey) {
return new FlussTrinoColumnHandle(
name, type, ordinalPosition, isPartitionKey, primaryKey,
isNullable, comment, defaultValue);
}

/**
* Create a new column handle with nullable flag
*/
public FlussTrinoColumnHandle withNullable(boolean nullable) {
return new FlussTrinoColumnHandle(
name, type, ordinalPosition, isPartitionKey, isPrimaryKey,
nullable, comment, defaultValue);
}

/**
* Create a new column handle with comment
*/
public FlussTrinoColumnHandle withComment(String commentText) {
return new FlussTrinoColumnHandle(
name, type, ordinalPosition, isPartitionKey, isPrimaryKey,
isNullable, Optional.ofNullable(commentText), defaultValue);
}

/**
* Create a new column handle with default value
*/
public FlussTrinoColumnHandle withDefaultValue(String defaultVal) {
return new FlussTrinoColumnHandle(
name, type, ordinalPosition, isPartitionKey, isPrimaryKey,
isNullable, comment, Optional.ofNullable(defaultVal));
}

/**
* Check if this column can be used for partition pruning
*/
public boolean canPrunePartition() {
return isPartitionKey;
}

/**
* Check if this column can be used for point queries
*/
public boolean canPointQuery() {
return isPrimaryKey;
}

/**
* Get column metadata for Trino
*/
public io.trino.spi.connector.ColumnMetadata toColumnMetadata() {
io.trino.spi.connector.ColumnMetadata.Builder builder =
io.trino.spi.connector.ColumnMetadata.builder()
.setName(name)
.setType(type)
.setNullable(isNullable);

comment.ifPresent(builder::setComment);

return builder.build();
}

@Override
public boolean equals(Object obj) {
if (this == obj) {
return true;
}
if (obj == null || getClass() != obj.getClass()) {
return false;
}
FlussTrinoColumnHandle that = (FlussTrinoColumnHandle) obj;
return ordinalPosition == that.ordinalPosition &&
isPartitionKey == that.isPartitionKey &&
isPrimaryKey == that.isPrimaryKey &&
isNullable == that.isNullable &&
Objects.equals(name, that.name) &&
Objects.equals(type, that.type) &&
Objects.equals(comment, that.comment) &&
Objects.equals(defaultValue, that.defaultValue);
}

@Override
public int hashCode() {
return Objects.hash(name, type, ordinalPosition, isPartitionKey,
isPrimaryKey, isNullable, comment, defaultValue);
}

@Override
public String toString() {
return "FlussTrinoColumnHandle{" +
"name='" + name + '\'' +
", type=" + type +
", ordinalPosition=" + ordinalPosition +
", isPartitionKey=" + isPartitionKey +
", isPrimaryKey=" + isPrimaryKey +
", isNullable=" + isNullable +
", comment=" + comment +
", defaultValue=" + defaultValue +
'}';
}
}
```
4. Enhanced Module Configuration
FlussTrinoMetadataModule.java

```
package com.alibaba.fluss.trino;

import com.google.inject.Binder;
import com.google.inject.Module;
import com.google.inject.Scopes;

/**
* Guice module for enhanced metadata management in Fluss Trino connector.
*/
public class FlussTrinoMetadataModule implements Module {

@Override
public void configure(Binder binder) {
// Bind enhanced metadata components
binder.bind(FlussTrinoMetadata.class).in(Scopes.SINGLETON);
binder.bind(FlussTrinoTableHandleResolver.class).in(Scopes.SINGLETON);
binder.bind(FlussTrinoColumnHandleResolver.class).in(Scopes.SINGLETON);

// Bind schema and table discovery services
binder.bind(FlussSchemaProvider.class).in(Scopes.SINGLETON);
binder.bind(FlussTableProvider.class).in(Scopes.SINGLETON);
}
}
```
5. Optimize module configuration
FlussTrinoOptimizationModule.java

```
package com.alibaba.fluss.trino;

import com.google.inject.Binder;
import com.google.inject.Module;
import com.google.inject.Scopes;

/**
* Guice module for query optimization features in Fluss Trino connector.
*/
public class FlussTrinoOptimizationModule implements Module {

@Override
public void configure(Binder binder) {
// Bind optimization components
binder.bind(FlussPredicatePushdown.class).in(Scopes.SINGLETON);
binder.bind(FlussColumnPruning.class).in(Scopes.SINGLETON);
binder.bind(FlussPartitionPruning.class).in(Scopes.SINGLETON);
binder.bind(FlussLimitPushdown.class).in(Scopes.SINGLETON);
binder.bind(FlussAggregatePushdown.class).in(Scopes.SINGLETON);

// Bind Union Read components
binder.bind(FlussUnionReadManager.class).in(Scopes.SINGLETON);
binder.bind(FlussLakehouseReader.class).in(Scopes.SINGLETON);
}
}
```
2. Data Reading Component
Refer to the source code implementation of Flink, and develop:

- FlussTrinoRecordSetProvider: Dataset Provider
- FlussTrinoSplit: Database Sharding Processing
- FlussTrinoSplitManager: Sharding Manager

3. Metadata Integration
Implementation using the existing Client connection interface:

- FlussTrinoMetadata: Trino metadata interface implementation
- FlussTrinoConnectorSplit: Connector Sharding
- Supports Trino's predicate pushdown and column pruning

Phase 3: Lakehouse Integration Enhancement

Based on the existing Lakehouse architecture, implement the Union Read feature:

1. Real-time data reading
- Read the latest data directly from TabletServer
- Supports streaming and batch query modes
- Implement the ConnectorPageSource interface of Trino

2. Historical Data Reading
- Read historical data through the Lakehouse storage layer
- Supports columnar storage formats such as Parquet/ORC
- Utilize Trino's existing file system connector

Phase 4: Query Optimization

1. Predicate Pushdown
Based on the optimization experience of Flink, achieve:
- Partition Cropping
- Column Cropping
- Filter Condition Pushdown

2. Aggregation Pushdown
Referring to Flink's aggregation pushdown implementation, it supports aggregation operations such as COUNT(*) to be executed at the storage layer.

Implementation Priority

High Priority
1. Basic Reading Function: Implement Trino's basic queries on Fluss's primary key table and log table
2. Metadata Integration: Enables Trino to discover and access Fluss tables via Catalog
3. Predicate Pushdown: Implements basic filter pushdown optimization

Medium Priority
1. Union Read: Enables joint query of real-time data and historical data
2. Aggregation Pushdown : Supports pushdown of simple aggregation operations
3. Partition Table Support: Fully supports querying and optimization of partition tables

Low Priority
1. Write Support : Write data to Fluss via Trino (if needed)
2. Advanced Optimization: Complex Query Optimization and Performance Tuning
3. Security Integration: Integration with Trino's Authentication and Authorization System

Compatibility Considerations

Referring to Flink's multi-version support model, create corresponding connector modules for different versions of Trino to ensure API compatibility.

### Solution

_No response_

### Anything else?

_No response_

### Willingness to contribute

- [ ] I'm willing to submit a PR!

Contributor guide

No contributing guide indexed for this repository

Research direction

Start by comparing the existing Flink integration with the proposed fluss-trino module structure, then review the fluss-trino/pom.xml, fluss-trino-common, and version-specific 435, 436, and 437 modules. The stated entry point is FlussConnectorFactory; completion would require the planned connector components, Trino version modules, dependencies, and tests to support the query flow described in the issue.

Written by the indexing model from the issue text.

Assessment

Tech stack
java
Domain
backend, data-engineering, distributed-systems
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
30/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.