Flink Connector
Overview
The Apache Gravitino Flink connector implements the Catalog Store to manage the catalogs under Gravitino. This capability allows users to perform federation queries, accessing data from various catalogs through a unified interface and consistent access control.
The connector is published as a version-specific runtime JAR for each supported Flink minor version. Internally, Gravitino uses shared connector logic with version-specific catalog and factory entry points so that each runtime artifact matches the Flink APIs it runs against.
Capabilities
- Supports Hive catalog
- Supports Iceberg catalog
- Supports Paimon catalog
- Supports Jdbc catalog
- Supports most DDL and DML SQLs.
Prerequisites
- Scala 2.12
- Flink 1.18, 1.19, or 1.20
- JDK 8, 11 or 17
Usage
- Build or download the Gravitino Flink connector runtime JAR that matches your Flink minor version, and place it in the classpath of Flink.
| Flink version | Runtime artifact |
|---|---|
| 1.18 | gravitino-flink-connector-runtime-1.18_2.12-${gravitino-version}.jar |
| 1.19 | gravitino-flink-connector-runtime-1.19_2.12-${gravitino-version}.jar |
| 1.20 | gravitino-flink-connector-runtime-1.20_2.12-${gravitino-version}.jar |
Do not mix runtime JARs from different Flink minor versions in the same Flink deployment.
- Configure the Flink configuration to use the Gravitino flink connector.
| Property | Type | Default Value | Description | Required | Since Version |
|---|---|---|---|---|---|
| table.catalog-store.kind | string | generic_in_memory | The Catalog Store name, it should set to gravitino. | Yes | 0.6.0-incubating |
| table.catalog-store.gravitino.gravitino.metalake | string | (none) | The metalake name that flink connector used to request to Gravitino. | Yes | 0.6.0-incubating |
| table.catalog-store.gravitino.gravitino.uri | string | (none) | The uri of Gravitino server address. | Yes | 0.6.0-incubating |
| table.catalog-store.gravitino.gravitino.enableSessionCatalogSupport | boolean | false | Whether to enable support for Flink's session catalog in the Gravitino catalog store. | No | 1.3.0 |
| table.catalog-store.gravitino.gravitino.client. | string | (none) | The configuration key prefix for the Gravitino client config. | No | 1.0.0 |
When table.catalog-store.gravitino.gravitino.enableSessionCatalogSupport is set to true, Gravitino uses GravitinoSessionCatalogStore, which combines a GravitinoCatalogStore (backed by the Gravitino server) with an in-memory store to support Flink's session catalog. When false (the default), only GravitinoCatalogStore is used.
When session catalog support is enabled, the following behaviors apply:
- CREATE CATALOG: Gravitino-managed catalogs (e.g.
gravitino-hive,gravitino-iceberg) are persisted to the Gravitino server; non-Gravitino-managed catalogs (e.g.hive,jdbc, or any custom connector type) are stored in the in-memory store only. - GET / USE CATALOG: The in-memory store is checked first. If the catalog is not found there, it is retrieved from the Gravitino server.
- DROP CATALOG: The in-memory store is checked first. If the catalog exists there it is removed from memory; otherwise it is removed from the Gravitino server.
- SHOW / LIST CATALOGS: Returns the combined set of catalogs from both the in-memory store and the Gravitino server.
- Session scope: Catalogs stored only in memory are session-scoped and will not survive when Flink restarts.
- Name conflict: If a catalog with the same name exists in both stores, the in-memory entry takes precedence.
To configure the Gravitino client, use properties prefixed with table.catalog-store.gravitino.gravitino.client.. These properties will be passed to the Gravitino client after removing the table.catalog-store.gravitino. prefix.
Example: Setting table.catalog-store.gravitino.gravitino.client.socketTimeoutMs is equivalent to setting gravitino.client.socketTimeoutMs for the Gravitino client.
Note: Invalid configuration properties will result in exceptions. Please see Gravitino Java client configurations for more support client configuration.
Set the flink configuration in flink-conf.yaml.
table.catalog-store.kind: gravitino
table.catalog-store.gravitino.gravitino.metalake: metalake_demo
table.catalog-store.gravitino.gravitino.uri: http://localhost:8090
table.catalog-store.gravitino.gravitino.enableSessionCatalogSupport: true
table.catalog-store.gravitino.gravitino.client.socketTimeoutMs: 60000
table.catalog-store.gravitino.gravitino.client.connectionTimeoutMs: 60000
Or you can set the flink configuration in the TableEnvironment.
final Configuration configuration = new Configuration();
configuration.setString("table.catalog-store.kind", "gravitino");
configuration.setString("table.catalog-store.gravitino.gravitino.metalake", "metalake_demo");
configuration.setString("table.catalog-store.gravitino.gravitino.uri", "http://localhost:8090");
configuration.setBoolean("table.catalog-store.gravitino.gravitino.enableSessionCatalogSupport", true);
configuration.setString("table.catalog-store.gravitino.gravitino.client.socketTimeoutMs", "60000");
configuration.setString("table.catalog-store.gravitino.gravitino.client.connectionTimeoutMs", "60000");
EnvironmentSettings.Builder builder = EnvironmentSettings.newInstance().withConfiguration(configuration);
TableEnvironment tableEnv = TableEnvironment.create(builder.inBatchMode().build());
- Add necessary jar files to Flink's classpath.
To run Flink with Gravitino connector and then access the data source like Hive, you may need to put additional jars to Flink's classpath. You can refer to the Flink document for more information.
- Execute the Flink SQL query.
Suppose there is only one hive catalog with the name catalog_hive in the metalake metalake_demo.
// use hive catalog
USE CATALOG catalog_hive;
CREATE DATABASE db;
USE db;
SET 'execution.runtime-mode' = 'batch';
SET 'sql-client.execution.result-mode' = 'tableau';
CREATE TABLE hive_students (id INT, name STRING);
INSERT INTO hive_students VALUES (1, 'Alice'), (2, 'Bob');
SELECT * FROM hive_students;
Catalog Naming Restrictions
When creating catalogs that will be used with the Flink connector, the catalog name cannot start with a number. This is a Flink limitation. For example:
- ✅ Valid:
catalog_hive,hive_catalog,my_catalog_1 - ❌ Invalid:
1_catalog,123catalog,2hive
If you create a catalog with a name starting with a number, it will not be accessible from Flink.
Datatype Mapping
Gravitino flink connector support the following datatype mapping between Flink and Gravitino.
| Flink Type | Gravitino Type | Since Version |
|---|---|---|
array | list | 0.6.0-incubating |
bigint | long | 0.6.0-incubating |
binary | fixed | 0.6.0-incubating |
boolean | boolean | 0.6.0-incubating |
char | char | 0.6.0-incubating |
date | date | 0.6.0-incubating |
decimal | decimal | 0.6.0-incubating |
double | double | 0.6.0-incubating |
float | float | 0.6.0-incubating |
integer | integer | 0.6.0-incubating |
map | map | 0.6.0-incubating |
null | null | 0.6.0-incubating |
row | struct | 0.6.0-incubating |
smallint | short | 0.6.0-incubating |
time | time | 0.6.0-incubating |
timestamp | timestamp without time zone | 0.6.0-incubating |
timestamp without time zone | timestamp without time zone | 0.6.0-incubating |
timestamp with time zone | timestamp with time zone | 0.6.0-incubating |
timestamp with local time zone | timestamp with time zone | 0.6.0-incubating |
timestamp_ltz | timestamp with time zone | 0.6.0-incubating |
tinyint | byte | 0.6.0-incubating |
varbinary | binary | 0.6.0-incubating |
varchar | string | 0.6.0-incubating |