Connect Apache Flink to LindormTable and use Lindorm tables as dimension tables or result tables in your Flink jobs. You can access these tables using Flink SQL or Flink DataStream.
Background information
You can use Lindorm LindormTable as a dimension table or result table in Flink, and access LindormTable through Flink SQL or Flink DataStream.
Choose a connection method
In LindormTable, tables are divided into two types based on how they are created: HBase tables (tables created and written to by using the HBase API) and SQL tables (tables created and written to by using Lindorm SQL). Accessing these two types of tables requires different interfaces. Therefore, before creating a Flink task, you need to determine the connector type and then determine the connection address of the LindormTable based on the connector type.
Determine the table type to access
In Lindorm, you can use Lindorm SQL to determine the type of LindormTable that a Flink task in development needs to access:
Use lindorm-cli (or Lindorm Insight or DMS directly) to connect to the Lindorm wide table engine.
Run the following SQL statement to check the
IS_HBASE_LIKEattribute of the table.If the attribute value is TRUE, the table is an HBase table.
If the attribute value is FALSE, the table is a SQL table.
SHOW TABLE VARIABLES FROM table_name LIKE 'IS_HBASE_LIKE';
For information about how to use
lindorm-clito connect to the wide table engine, see Connect to and use the wide table engine with Lindorm-cli.For the detailed syntax of
SHOW TABLE VARIABLE, see SHOW VARIABLES.
Choose the connector type
Based on the Flink product form you have selected, decide which connector to use to access the LindormTable.
Table type | Community Flink | Realtime Compute for Apache Flink |
HBase table | (Supported as dimension table and result table) | Cloud-native multi-model database Lindorm connector (Supported as dimension table and result table) |
SQL table | (Supported as result table) | Cloud-native multi-model database Lindorm connector (Supported as dimension table and result table) |
Realtime Compute for Apache Flink is a managed Flink service on Alibaba Cloud. For more information, see Realtime Compute for Apache Flink. Note that a Flink cluster built with open source Flink on Alibaba Cloud ECS is still classified as community Flink in the table above.
Get the connection information of LindormTable
Scenario 1: Use the open source HBase connector or the cloud-native multi-model database Lindorm connector
In this scenario, the connection address must be the HBase Java API access address (VPC) of the wide table engine.
On the Lindorm instance details page, click Database connections on the left-side menu, and select the Wide Table Engine tab. In the Connect through HBase-compatible address area, obtain the VPC address in the format ofld-<instance-ID>-proxy-lindorm.lindorm.rds.aliyuncs.com:30020. In the Connect through MySQL-compatible address area, obtain the MySQL-compatible address in the format ofld-<instance-ID>-proxy-sql-lindorm.lindorm.rds.aliyuncs.com:33060.Scenario 2: Use the JDBC connector
In this scenario, the connection address must be the MySQL-compatible address (VPC) of the wide table engine.
On the Lindorm instance details page, click Database connections in the left-side navigation pane, and select the Wide Table Engine tab. In the Connect through HBase-compatible address area, you can view the VPC address of the HBase Java API (in the format ofld-<instance-ID>-proxy-lindorm.rds.aliyuncs.com:30020).
If the Flink task uses a newly created Lindorm user to access the LindormTable, ensure that the user has read and write permissions on the Flink tables. For information about how to grant permissions, see Grant permissions to specified users.
For detailed descriptions of the various connection addresses of Lindorm wide tables, see View the connection addresses of the wide table engine.
Method for a Flink job to access LindormTable
You can develop with the general approach of the selected real-time compute framework. To access the LindormTable in a job, refer to the following documents based on the connector you have selected to enable your compute job to access the LindormTable.
Prerequisites
Develop a job that uses community Flink to access the LindormTable
To use the open source HBase connector to access the LindormTable, ensure that the wide table engine is version 2.4.3 or later.
To use the open source JDBC connector to access the LindormTable, ensure that the wide table engine is version 2.6.5.2 or later, and that you have Enable the MySQL compatibility feature.
For information about how to view or upgrade the current version, see LindormTable release notes and Upgrade the minor engine version of a Lindorm instance.
If you use Realtime Compute for Apache Flink to develop a job to access the LindormTable, there is no version restriction on the wide table engine.
Ensure that the environment where the Flink cluster is located has network connectivity with the Lindorm instance, and that the client IP addresses are added to the Lindorm whitelist. For information about how to add IP addresses to the whitelist, see Configure a whitelist.
Job development with community Flink
Open source HBase connector
When you use community Flink to develop a job to access HBase tables, if you want to access the tables over the Internet, or if the target Lindorm instance is a Lindorm single-node instance, you must upgrade the SDK and change the configuration before performing subsequent operations.
For details, see Step 1 in Connect to and use LindormTable using the HBase Java API.
For information about how to use the open source HBase connector to create a dimension table and a result table, see HBase connector documentation.
Open source JDBC connector
When you use the open source JDBC connector to access a LindormTable, currently only using the LindormTable as a result table is supported. For overall usage, refer to the JDBC connector official documentation. However, there are some points that require special attention, as follows:
Dependency requirements
Compared with the relatively broad range of dependency package versions in the official documentation, the dependency package versions used by the JDBC connector to access LindormTables are currently restricted to the following list:flink-connector-jdbc-core-4.0.0-2.0.jarflink-connector-jdbc-mysql-4.0.0-2.0.jarmysql-connector-j-8.3.0.jar
The dependency package of the MySQL JDBC driver can be downloaded from the community.
JDBC connector parameters
Because using as a source table and dimension table is not supported, parameters related to source table and dimension table functionality are not supported (such as parameters prefixed with scan likescan.fetch-size, and parameters prefixed withlookuplikelookup.cache).
Recommendations for some JDBC connector parameters:url: We recommend that you follow Use Java JDBC APIs to develop applications for configuration.
username: Use the username created in the Lindorm instance.
password: The password of the username.
connector, table-name: Follow the JDBC connector community recommendations.
sink parameters: Fine-tune based on the actual situation of the job.
Data type mapping
The mapping between Flink data types and Lindorm data types can generally follow that of MySQL data types (refer to the Data type mapping section of the JDBC connector). However, some data types in Lindorm do not align with MySQL. For example, the following MySQL types claimed to be supported in the JDBC connector are not supported by Lindorm:MEDIUMINT type
DATETIME type
UNSIGNED types other than BIGINT UNSIGNED
The following example uses Flink SQL to define a job that accesses a LindormTable via the JDBC connector. In this example, assume that a table named testflink has been defined in the Lindorm wide table engine.
# Create Flink table and start the job
CREATE TABLE source_table(
c1 INT,
c2 STRING
) WITH (
'connector' = 'datagen',
'rows-per-second' = '2',
'fields.c2.length' = '5',
'fields.c1.min' = '1',
'fields.c1.max' = '100'
);
CREATE TABLE sink_table(
c1 INT,
c2 STRING
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://ld-xxxxx-proxy-lindorm.lindorm.rds.aliyuncs.com:33060/default?sslMode=disabled&allowPublicKeyRetrieval=true&useServerPrepStmts=true&useLocalSessionState=true&rewriteBatchedStatements=true&cachePrepStmts=true&prepStmtCacheSize=300&prepStmtCacheSqlLimit=50000000',
'username' = 'root',
'password' = 'root',
'table-name' = 'testflink'
);
INSERT INTO sink_table SELECT * FROM source_table;Job development with Realtime Compute for Apache Flink
Cloud-native multi-model database Lindorm connector
In Realtime Compute for Apache Flink job development, you can develop jobs that access the Lindorm wide table engine by using Flink SQL. For details about Realtime Compute for Apache Flink job development, see Job development overview.
For information about how to use the connector to create a dimension table and a result table, see Cloud-native multi-model database Lindorm connector documentation.