Flink SQL Streaming nodes in DataWorks Data Studio let you define real-time processing logic in standard SQL. They support robust state management, fault tolerance, event-time and processing-time semantics, and integrate with systems such as Kafka and HDFS. This topic describes how to develop, configure, and run a Flink SQL Streaming node to process real-time data.
Prerequisites
-
A compute resource for Realtime Compute for Apache Flink is associated in Administration. For more information, see Bind a compute engine.
-
A Flink SQL Streaming node is created. For more information, see Create a node for a scheduling workflow.
-
You have granted the required OpenAPI permissions to the RAM user or RAM role that DataWorks uses to call Realtime Compute for Apache Flink APIs. These permissions allow DataWorks to submit and deploy node tasks to a Flink cluster.
{ "Version": "1", "Statement": [ { "Effect": "Allow", "Action": ["stream:CreateDeployment", "stream:UpdateDeployment", "stream:GetDeployment", "stream:DeleteDeployment"], "Resource": ["*"] } ] }
Limitations
-
This node cannot be used in a workflow; it must be developed and run as a standalone node.
-
Only serverless resource groups are supported. Legacy exclusive resource groups for scheduling are not supported.
Step 1: Develop the Flink SQL Streaming node
On the Flink SQL Streaming node editing page, develop the node task.
Develop SQL code
In the SQL editor, you can define variables using the ${variable_name} format. Assign values to these variables in the Script Parameters section of the Real-Time configuration panel to pass parameters dynamically in scheduling scenarios. For example:
--Create the source table datagen_source.
CREATE TEMPORARY TABLE datagen_source(
name VARCHAR
) WITH (
'connector' = 'datagen'
);
--Create the result table blackhole_sink.
CREATE TEMPORARY TABLE blackhole_sink(
name VARCHAR
) WITH (
'connector' = 'blackhole'
);
--Insert data from the source table into the result table.
INSERT INTO blackhole_sink
SELECT
name
FROM datagen_source WHERE LENGTH(name) > ${name_length};
In this example, the value of the parameter name_length is 5. This parameter filters the data to process only records where the name length is greater than 5 characters.
Step 2: Configure the Flink SQL Streaming node
Configure the following parameters for the Flink SQL Streaming node based on your business requirements.
Configure Flink resources
In the Flink resource information section of the Real-Time configuration panel, configure the following parameters based on the selected Resource Mode. For more information, see Configure Flink resources.
|
Parameter |
Description |
|
Flink cluster |
The fully managed Flink compute resource associated in Administration. |
|
Flink engine version |
The engine version to use. Select a version based on your needs. |
|
Resource Group |
Select a serverless resource group that has network connectivity with Flink. |
|
Resource Mode supports the following two modes. For more information, see Configure Flink resources.
Configure parameters based on the resource mode you selected. Understanding the Flink architecture helps you configure parameters more effectively. For more information, see Flink Architecture | Apache Flink. |
|
|
Basic mode |
|
|
Job Manager CPU |
JobManager requires at least 0.5 CPU cores and 2 GiB of memory for stable operation. The recommended configuration is 1 CPU core and 4 GiB of memory, with a maximum of 16 CPU cores. Adjust based on cluster scale and job complexity. |
|
Job Manager Memory |
JobManager memory affects scheduling and management capacity. The recommended range is 2 GiB to 64 GiB. Adjust based on cluster scale and job requirements. |
|
Task Manager CPU |
TaskManager CPU affects task processing capability. At least 0.5 CPU cores and 2 GiB of memory are recommended, with a preferred configuration of 1 CPU core and 4 GiB of memory. The maximum is 16 CPU cores. Adjust based on your requirements. |
|
Task Manager Memory |
TaskManager memory determines the data volume and processing performance. The memory size must be at least 2 GiB and can be set up to 64 GiB. |
|
Concurrency |
The number of parallel task executions in a Flink job. Higher concurrency improves processing speed and resource utilization. Set this value based on cluster resources and job characteristics. |
|
Number of slots per TaskManager |
The number of slots per TaskManager, which determines how many tasks it can run in parallel. Adjust slots to optimize resource utilization and parallel processing. |
|
Expert mode |
|
|
Job Manager CPU |
JobManager requires at least 0.25 CPU cores and 1 GiB of memory for stable operation, with a maximum of 16 CPU cores. Adjust based on cluster scale and job complexity. |
|
Job Manager Memory |
JobManager memory affects scheduling and management capacity. The recommended range is 1 GiB to 64 GiB. Adjust based on cluster scale and job requirements. |
|
Number of slots per TaskManager |
The number of slots per TaskManager, which determines how many tasks it can run in parallel. Adjust slots to optimize resource utilization and parallel processing. |
|
Multiple SSG mode |
By default, all operators share a single slot sharing group, so you cannot configure resources for each operator individually. Enable Multiple SSG mode to assign each operator its own independent slot, then configure resources on the corresponding slot. |
(Optional) Configure script parameters
In the Script Parameters section of the Real-Time configuration panel in the right navigation pane, click Add parameters and edit the Parameter name and Parameter Value to dynamically use them in your code.
(Optional) Configure Flink running parameters
In the Flink running parameters section of the Real-Time configuration panel in the right navigation pane, configure the following parameters. For more information, see Configure Flink running parameters.
|
Parameter |
Description |
|
System Checkpoint Interval |
The time interval at which Flink performs periodic checkpoints. A shorter interval reduces fault recovery time but increases system overhead. If left empty, checkpoints are disabled. |
|
Minimum time interval between two system checkpoints |
The minimum wait time between consecutive checkpoints, preventing frequent checkpoints from affecting performance. This ensures a minimum gap between two checkpoints when the maximum parallelism for checkpoints is 1. |
|
State Data Expiration Time |
The maximum time that state data can be retained without being accessed or updated. The default is 36 hours, after which state information automatically expires and is cleared to optimize storage and resource usage. Important
This default is based on cloud best practices and differs from the open-source default (0, meaning state never expires). |
|
Others |
Additional Flink running parameters. For example: Note
For more information about parameter configurations, see Configure Flink running parameters. |
After the task is configured, click Save to save the node task.
Step 3: (Optional) Debug the Flink SQL Streaming node
Before deploying the node to production, you can debug it by running the node code against uploaded mock data. This lets you verify SQL logic and upstream/downstream data without deploying the task to Operation Center.
The debug feature is available through an allowlist. To use this feature, submit a ticket to request access.
Configure Flink resource information
In the Flink resource information section of the Run Configuration panel on the right side of the node editing page, configure the parameters as described in the following table.
|
Parameter |
Description |
|
Flink Debug Cluster |
The Flink session cluster for running the debug task. Required. The drop-down list shows existing session clusters under the current compute resource and their running status. Only clusters in the Running state can be selected. If no clusters are available in the list, click Create Cluster to go to the Realtime Compute for Apache Flink console and create a session cluster. |
|
Flink Engine Version |
The Flink engine version of the selected session cluster. Automatically displayed based on the cluster selection. |
|
Timeout |
The maximum runtime for a single debug task, in minutes. Default: 30 minutes. The debug task stops automatically after this duration. |
After you switch the compute resource for the current node, the selected Flink Debug Cluster and uploaded debug data are cleared. You must reselect a cluster and re-upload the data.
Prepare debug data
In the Debug Data section of the Run Configuration panel, prepare mock data for the source tables referenced in your code.
-
Click Generate Template. The system parses the source tables referenced in the current SQL and generates corresponding table name records in the list below. Previously uploaded data is not cleared.
-
In the Actions column of a table name record, click Download Template to download a CSV template that matches the schema of the source table.
-
Fill in the debug data locally in the column order of the template and save the file in CSV format.
-
In the Actions column of a table name record, click Upload and select the completed CSV file to upload. After the upload succeeds, the Status column displays Enabled.
-
(Optional) After the upload succeeds, you can click Preview to view the data content in the bottom panel. To modify the data, re-upload a CSV file to overwrite the existing data.
-
If you do not want the mock data of a specific source table to participate in the current debug session, click Disable to change the status to Disabled. To re-enable it, click Enable. Only data with the Enabled status is used in the debug session.
Before uploading debug data, select a Flink Debug Cluster first. Otherwise, you are prompted to select a compute resource first.
Debug data supports only the CSV format, with a maximum file size of 1 MB. The first row of the CSV file must contain column names, and UTF-8 encoding is recommended.
Run the debug task
After the debug data is prepared, click the Run button on the editor toolbar (or press F8). The system submits the code, mock data, and Flink resource information to the selected session cluster for execution.
If your code uses parameters in the ${variable_name} format, make sure that you have assigned values to the variables in the Script Parameters section. During debugging, the system replaces the placeholders in the code with the assigned values before submission.
View debug results
After the debug task runs, the result area at the bottom of the node shows the following information:
-
Code: The SQL code submitted to the Flink engine (with variable substitution applied).
-
Logs: Runtime logs and error information.
-
Query results: The output data from the debug task.
Step 4: Start the Flink SQL Streaming node
-
Deploy the Flink SQL Streaming node.
The task must be deployed to Operation Center before it can run. Follow the on-screen instructions to deploy the node. For more information, see Deploy a node.
NoteThis operation also deploys the task to the Flink VVP workspace. You can view tasks deployed through DataWorks in Flink VVP Operation Center > Job O&M.
-
Start the Flink SQL Streaming node.
After the task is deployed, click Go to operation and maintenance below Deploy to production environment. In Operation Center, go to , find the task that you want to start, and click Start in the Operation column to start the real-time task and view its running status.