After you configure a MaxCompute catalog, you can directly access tables stored in MaxCompute in Realtime Compute for Apache Flink jobs, without defining a schema. This topic explains how to create, view, use, and delete MaxCompute catalogs in Realtime Compute for Apache Flink.
Background
The MaxCompute Catalog queries the MaxCompute service to retrieve the schema of physical tables. This allows you to retrieve detailed field information in Flink SQL without declaring a table schema. The MaxCompute Catalog provides the following features:
-
In a MaxCompute Catalog, the database name corresponds to the project name in MaxCompute. You can switch between databases to use tables from different MaxCompute projects.
-
In a MaxCompute Catalog, table names correspond to the names of physical tables in MaxCompute, and data types are mapped automatically. This eliminates the need to manually register MaxCompute tables with DDL statements, thereby improving development efficiency and accuracy.
-
You can directly use tables from a MaxCompute Catalog as source tables, dimension tables, and result tables in Flink SQL jobs.
-
Creating a table in a MaxCompute Catalog automatically creates a corresponding physical table in the MaxCompute service and maps the data types, thereby improving development efficiency.
This topic describes how to manage a MaxCompute Catalog:
Limitations
-
The MaxCompute catalog requires Ververica Runtime (VVR) 6.0.7 or later.
-
Creating databases, which correspond to projects in MaxCompute, is not supported.
-
Modifying table schemas is not supported.
Create a MaxCompute catalog
You can create a MaxCompute catalog by using the Console UI or Flink SQL. We recommend using the Console UI.
Console
-
Go to the Catalogs page.
-
Log on to the Realtime Compute for Apache Flink console. On the Workspaces page, find the target workspace and click Console in the Actions column.
-
In the left-side navigation pane, click Catalogs.
-
-
Click Create Catalog, select ODPS, and then click Next.
-
Configure the parameters.
ImportantYou cannot change the following configuration parameters after the catalog is created. To make changes, you must delete and recreate the catalog.
Parameter
Description
Type
Required
Notes
Catalog name
The name of the MaxCompute catalog.
String
Yes
A user-defined name in English.
endpoint
The endpoint for connecting to the MaxCompute service.
String
Yes
For a list of endpoints, see Endpoints.
accessId
The AccessKey ID of the Alibaba Cloud account used to access the MaxCompute service.
String
Yes
The account requires administrative permissions on the projects that the catalog will access.
accessKey
The AccessKey Secret of the Alibaba Cloud account used to access the MaxCompute service.
String
Yes
None.
Project
The MaxCompute project to use as the default database in the catalog.
String
No
If this parameter is not set, the default project is
default.NoteAfter the catalog is created, the metadata displays the project you specified and all other projects created by the provided Alibaba Cloud account.
catalog.schema.enabled
Specifies whether a database in a Flink catalog maps to a schema in MaxCompute.
Boolean
No
A schema is a mechanism for organizing tables, resources, and UDFs within a project. A project can contain multiple schemas. For more information, see Schema operations.
Valid values:
-
false (default): Maps a database in a Flink catalog to a MaxCompute project. Use this option for MaxCompute services where the schema feature is disabled.
-
true: Maps a database in a Flink catalog to a MaxCompute schema. Use this option for MaxCompute services where the schema feature is enabled.
-
-
Click OK.
Once created, the catalog is listed on the Catalogs page.
ImportantIf the AccessKey pair used to create the catalog does not have permissions for a specific project, information about that project will not be displayed in the metadata. However, this does not affect normal read and write operations on the catalog.
Flink SQL
-
In the SQL Editor, enter the command to create a MaxCompute catalog.
CREATE CATALOG `<catalogName>` WITH ( 'type' = 'odps', 'endpoint' = '<odpsEndpoint>', 'accessId' = '<aliyunAccountAccessId>', 'accessKey' = '<aliyunAccountAccessKey>', 'project' = '<defaultProject>', 'userAccount' = '<RAMUserAccount>' );The following table describes the parameters.
Parameter
Description
Type
Required
Notes
catalogName
The name of the MaxCompute catalog.
String
Yes
A user-defined name in English.
type
The catalog type.
String
Yes
The value must be
odps.endpoint
The endpoint for connecting to the MaxCompute service.
String
Yes
For a list of endpoints, see Endpoints.
accessId
The AccessKey ID of the Alibaba Cloud account used to access the MaxCompute service.
String
Yes
The account requires administrative permissions on the projects that the catalog will access.
accessKey
The AccessKey Secret of the Alibaba Cloud account used to access the MaxCompute service.
String
Yes
None.
project
The MaxCompute project to use as the default database in the catalog.
String
No
If this parameter is not set, the default project is
default.userAccount
The name of the Alibaba Cloud account or RAM user.
String
No
If you are using an AccessKey pair from a RAM user that has admin permissions only on specific projects, you must set this parameter to the account name. Example:
RAM$[<account_name>:]<RAM_name>. The MaxCompute catalog will then list only the projects that the account can access.For more information about managing user permissions in MaxCompute, see User planning and management.
-
Select the statement and click Run in the left margin.
CREATE CATALOG `<catalogName>` WITH ( 'type' = 'odps', 'endpoint' = '<odpsEndpoint>', 'accessId' = '<aliyunAccountAccessId>', 'accessKey' = '<aliyunAccountAccessKey>', 'project' = '<defaultProject>' );
View MaxCompute Catalog
Console UI (recommended)
-
Navigate to the Catalogs page.
-
Log on to the Realtime Compute for Apache Flink console.
-
Find the target workspace and click Console in the Actions column.
-
Click Catalogs.
-
-
On the Catalog List page, check the Name and Type columns.
To view the databases and tables, click View.
Flink SQL
-
In the SQL Editor, enter the following command:
DESCRIBE `<catalogName>`.`<projectName>`.`<tableName>`;Parameter
Description
catalogName
The name of the MaxCompute Catalog.
projectName
The name of the MaxCompute project.
tableName
The name of the physical MaxCompute table.
-
Select the statement and click Run in the line-number area on the left.
After the statement runs successfully, the table schema appears in the Results tab below the editor.
Use MaxCompute Catalog
Create MaxCompute physical tables
When you create a table in a MaxCompute Catalog using a Flink SQL DDL statement, a corresponding physical table is automatically created in the specified MaxCompute project. Flink data types are automatically converted to their corresponding MaxCompute types. You can create both non-partitioned tables and partitioned tables.
Example of creating a non-partitioned table:
CREATE TABLE `<catalogName>`.`<projectName>`.`<tableName>` (
f0 INT,
f1 BIGINT,
f2 DOUBLE,
f3 STRING
);
After the statement executes, you can view the tables in the corresponding MaxCompute project. A non-partitioned table with the specified name is created. Its column names and data types match those defined in the Flink DDL statement.
Example of creating a partitioned table:
CREATE TABLE `<catalogName>`.`<projectName>`.`<tableName>` (
f0 INT,
f1 BIGINT,
f2 DOUBLE,
f3 STRING,
ds STRING
) PARTITIONED BY (ds);
To create a partitioned table, add the partition columns at the end of the schema in the Flink DDL and declare them in the PARTITIONED BY clause. After the statement executes, you can view the table in the corresponding MaxCompute project. A partitioned table with the specified name is created. In this example, f0, f1, f2, and f3 are regular columns, and ds is the partition column.
Column names in MaxCompute are all lowercase, whereas column names in Flink are case-sensitive. If a column name in a DDL statement contains uppercase letters, it is automatically converted to lowercase. If the DDL statement contains multiple columns that have the same name after being converted to lowercase, an error occurs.
Read data from MaxCompute tables
MaxCompute Catalog can retrieve the schema of a physical table from the MaxCompute service. This allows you to read data directly without declaring the table schema in Flink. For example:
SELECT * FROM `<catalogName>`.`<projectName>`.`<tableName>`;
By default, if you do not specify any parameters, Flink reads all data from all partitions. If you need to read from specific partitions or use the incremental source mode, you can specify options in a SQL comment, as described in the MaxCompute documentation. For example:
Read from a specific partition:
SELECT * FROM `<catalogName>`.`<projectName>`.`<tableName>`
/*+ OPTIONS('partition' = 'ds=230613') */;
Use the incremental source mode:
SELECT * FROM `<catalogName>`.`<projectName>`.`<tableName>`
/*+ OPTIONS('startPartition' = 'ds=230613') */;
Use a dimension table:
SELECT * FROM `<anotherTable>` AS l LEFT JOIN
`<catalogName>`.`<projectName>`.`<tableName>`
/*+ OPTIONS('partition' = 'max_pt()', 'cache' = 'ALL') */
FOR SYSTEM_TIME AS OF l.proc_time AS r
ON l.id = r.id;
You can set other source table and dimension table parameters supported by MaxCompute in the same way. Note that MaxCompute Catalog does not store watermark information. If you need to specify a watermark when reading data from a source table, use a CREATE TABLE ... LIKE ... statement. For example:
CREATE TABLE `<newTable>` ( WATERMARK FOR ts AS ts )
LIKE `<catalogName>`.`<projectName>`.`<tableName>`;
In this example, ts is a DATETIME column in the MaxCompute physical table. In Flink, you can use this column as the event time and add a watermark. After the table is created, all data read from newTable contains watermarks.
Write data to MaxCompute tables
MaxCompute Catalog supports writing data in static or dynamic partition mode. For more information, see MaxCompute. For example, if a MaxCompute physical table has two partition levels, ds and hh, you can use the following statements to write data:
-- Write to a static partition
INSERT INTO `<catalogName>`.`<projectName>`.`<tableName>`
/*+ OPTIONS('partition' = 'ds=20231024,hh=09') */
SELECT <otherColumns>, '20231024', '09' FROM `<anotherTable>`;
-- Write to a dynamic partition
INSERT INTO `<catalogName>`.`<projectName>`.`<tableName>`
/*+ OPTIONS('partition' = 'ds,hh') */
SELECT <otherColumns>, ds, hh FROM `<anotherTable>`;
In the SELECT statement, place the partition columns after the regular columns in the same order as their partition levels.
Delete a MaxCompute catalog
Deleting a MaxCompute catalog does not affect running deployments. However, deployments that use tables from the catalog fail if they are published or restarted. Proceed with caution.
Console
-
Go to the Catalogs page.
-
Log in to the Realtime Compute for Apache Flink console.
-
Find the target workspace and click Console in the actions column.
-
Click Catalogs.
-
-
On the Catalog list page, find the target catalog and click Delete in the actions column.
-
In the confirmation dialog, click Delete.
NoteAfter the catalog is deleted, verify that it no longer appears in the Catalogs navigation pane on the left.
Flink SQL
-
In the SQL editor, enter the following statement.
DROP CATALOG `<catalogName>`;Replace <catalogName> with the name of the MaxCompute catalog that you want to delete.
WarningThis operation does not affect running deployments. However, deployments that use tables from the catalog fail if they are published or restarted. Proceed with caution.
-
Select the statement, right-click, and select Run.
-
In the Catalogs navigation pane on the left, verify that the catalog no longer appears.
MaxCompute and Flink type mapping
For the data types that MaxCompute supports, see Data types (Version 2.0).
MaxCompute to Flink
When you read from a MaxCompute physical table, its field data types are mapped to Flink data types as shown in the following table.
|
MaxCompute type |
Flink type |
|
BOOLEAN |
BOOLEAN |
|
TINYINT |
TINYINT |
|
SMALLINT |
SMALLINT |
|
INT |
INTEGER |
|
BIGINT |
BIGINT |
|
FLOAT |
FLOAT |
|
DOUBLE |
DOUBLE |
|
DECIMAL(precision, scale) |
DECIMAL(precision, scale) |
|
CHAR(n) |
CHAR(n) |
|
VARCHAR(n) |
VARCHAR(n) |
|
STRING |
STRING |
|
BINARY |
BYTES |
|
DATE |
DATE |
|
DATETIME |
TIMESTAMP(3) |
|
TIMESTAMP |
TIMESTAMP(9) |
|
ARRAY |
ARRAY |
|
MAP |
MAP |
|
STRUCT |
ROW |
|
JSON |
STRING |
Flink to MaxCompute
When you use Flink DDL to create a MaxCompute table in a catalog, the Flink data types for the fields are mapped to MaxCompute data types as shown in the following table.
|
Flink type |
MaxCompute type |
|
BOOLEAN |
BOOLEAN |
|
TINYINT |
TINYINT |
|
SMALLINT |
SMALLINT |
|
INTEGER |
INT |
|
BIGINT |
BIGINT |
|
FLOAT |
FLOAT |
|
DOUBLE |
DOUBLE |
|
DECIMAL(precision, scale) |
DECIMAL(precision, scale) |
|
CHAR(n) |
CHAR(n) |
|
VARCHAR / STRING |
STRING |
|
BINARY |
BINARY |
|
VARBINARY / BYTES |
BINARY |
|
DATE |
DATE |
|
TIMESTAMP(n<=3) |
DATETIME |
|
TIMESTAMP(3<n<=9) |
TIMESTAMP |
|
ARRAY |
ARRAY |
|
MAP |
MAP |
|
ROW |
STRUCT |