Real-time tasks
By deeply integrating the powerful stream computing capabilities of Apache Flink, the AE DataOps Platform provides a higher-performance streaming data sync solution and supports real-time analysis and exploration of streaming data, building a future-ready enterprise data development platform with unified stream and batch processing, high real-time performance, and low operations costs.
1. Create a real-time task
1.1 Entry point
In the Dev module of the DataOps Platform, click + Realtime Task in the upper-right corner of the Realtime Tasks tab to open the creation page.
1.2 Fill in basic information
Real-time tasks support dual environments (development environment / production environment). The basic information is shared by both environments.
| Field | Description | Remarks |
|---|---|---|
| Name | Unique within the workspace, up to 48 characters; supports Chinese characters, English letters, digits, and underscores | Required |
| Flink version | The Flink engine version bound to the real-time task, selected from the drop-down list | Can't be changed after creation |
| Owner | The owner of the real-time task. The creator is the owner by default | Required |
| Remarks | Task description, up to 200 characters | Optional |
1.3 Development environment edit page
After the task is created, the Dev environment edit page opens. It mainly includes the following areas:
1.3.1 Code editing area
-
Supports Flink SQL syntax highlighting and autocomplete
-
Supports parameterized configuration. Reference parameters in the
${参数名}format -
The editor automatically parses the parameters in the code and shows them in the Used Parameter area below
-
The editor also parses table lineage, showing information such as the sources, destinations, and connector types of the tables involved in the Flink SQL
1.3.2 Task parameter management (sidebar)
Real-time tasks support the following two parameter types:
| Parameter type | Description |
|---|---|
| Workspace Params | Supports preset workspace parameters and custom workspace parameters |
| Task Params | Isolated at the real-time task level, with two types: "Plain Text" and "Expression" |
Ways to add parameters:
-
Method 1: Create manually — add parameters manually in the sidebar. When the task is submitted, the parameter values are parsed and passed to the execution engine
-
Method 2: Parse from code — define parameters directly in Flink SQL in the
${abc}format, and the system automatically parses them and shows them in the sidebar
Dynamic parameters:
When you need different parameter values in different environments, use dynamic parameters to specify the value for the current environment at runtime.
${env.dev: development environment value, env.product: production environment value}
Parameter validation rules (validated only when debugging/running, not when saving or releasing):
- Task Params: if
${abc}is declared in the SQL but has no value defined, debugging/running is blocked and you're asked to assign a value
- Workspace Params: if a workspace parameter referenced in the SQL has been deleted, debugging/running is blocked with a message that the parameter doesn't exist
2. Debug a real-time task
2.1 About debugging
Debugging verifies that a real-time task works correctly in the Dev environment, so you can test the task logic without releasing it to production.
When you click Debug, you can choose:
- Save only: triggers only the save validation process
- Save and debug: triggers the save validation + debug validation processes
Debugging isn't a required step before releasing. After saving, you can release directly without debugging first.
2.2 Validation process
Save validation:
-
SQL syntax validation
-
Table permission validation (read and write permissions on the source and target tables)
-
Parameter parsing and recording
Debug/run validation (in addition to save validation):
- Parameter completeness validation — all
${XX}parameters must have values - Workspace parameter existence validation
- Execution channel status validation
2.3 Debug run
After debugging starts successfully, the task runs in the execution channel of the development environment. You can:
- Go to the Flink UI to view the task's running status
- View real-time logs
- View runtime metrics
- Manually stop the debugging task
3. Save and release a real-time task
3.1 Save
Releasing (bringing online) a real-time task is one-way: a "Published" real-time task can't go back to the "Unpublished" status (unlike offline Flows, which can switch both ways).
Saving a task that's "Pending Publishing":
| Item | Description |
|---|---|
| Environment | Saved only in the Dev environment |
| Status | Pending Publishing |
| Version | The development environment shows the latest version; none in production |
| Available operations | Edit, Release, Delete |
Saving a task that's "Published":
| Item | Description |
|---|---|
| Environment | The new version is saved only in the Dev environment and must be released to production again |
| Status | Published (with changes) |
| Version | The development environment shows the latest version; production still has the previously released version |
| Available operations | Edit, Re-release, Run, Delete |
3.2 Save and release
Releasing syncs the real-time task version in the development environment to the production environment.
First release:
-
The version is updated in both the development and production environments
-
After a successful release, the status changes to "Published". Click Start to run the first instance
-
If the release fails, the status stays "Pending Publishing"
Re-releasing a published task:
-
The versions in both the development and production environments are updated
-
The status stays "Published"
-
After a successful release, click Start to run an instance with the new version
3.3 Release validation
The following validations run during release:
- Unsaved tasks can't be released (the release button is grayed out)
- Production environment table permission validation
- Parameter value completeness validation
- Table schema consistency validation
4. Run (start) a real-time task
4.1 Production environment detail page
After a task is released to production, you can view and manage it on the production environment detail page. The detail page contains the following modules:
| Module | Description |
|---|---|
| Real-time task basic information | Task name, owner, engine version, change status, creation/modification time, description |
| Operation setting | Execution channel selection (resource configuration), runtime configuration parameters |
| Instance content | SQL content, runtime configuration, table lineage, and task parameter values of the latest instance |
| Operation record | Event messages for the full lifecycle of instances |
| Log | Includes run logs and exception logs |
| Snapshot | Total number of Checkpoints/Savepoints, success/failure statistics, and history details |
4.2 Runtime parameter configuration
Before you Start a real-time task, complete the following configuration:
4.2.1 Parameter assignment
Code-parsed parameters need values assigned at startup, and the parameter values are recorded at the instance level.
4.2.2 Operation settings
Runtime resource configuration
| Config items | Description |
|---|---|
Channel (resource pool) selection | Channel name: the system automatically filters the available channels based on the task's Flink engine version
|
| Calculate mode |
|
| Parallelism | Resource specs are defined at the channel level. When an instance is submitted for execution, you choose the parallelism.
|
Runtime parameter configuration
| Category | Parameter | Description | Default value |
|---|---|---|---|
| Checkpoint settings | execution.checkpointing.interval | Checkpoint interval | 180 seconds |
execution.checkpointing.timeout | Checkpoint timeout | 3 minutes | |
execution.checkpointing.min-pause | Minimum pause between checkpoints | 60 seconds | |
Status data time out | table.exec.state.ttl | The expiration time of state data; 0 means it never expires When data first enters the system and is processed, it's stored in state memory. When the next record with the same primary key arrives, the system uses the previously stored state data for computation and updates its access time. This process is central to real-time computing because it depends on a continuous flow of data. If the data isn't accessed again within the configured TTL window, the system treats it as expired and clears it from the state store. Setting a reasonable TTL value keeps computation accurate and clears stale data promptly, effectively reducing state memory usage, which lowers the system's memory load and improves computing efficiency and system stability. | 0 |
| Restart strategy | restart-strategy.type | Restart strategy type Only when no restart strategy is configured does Flink decide whether to restart the job based on whether system checkpoints are enabled (if system checkpoints are enabled, the job restarts at the fixed interval set by Fixed Delay; if not, the job doesn't restart). If a Flink restart strategy is configured, the job restarts according to the configured strategy. The parameter values are as follows:
| Fixed Delay |
4.2.3 Startup strategy
When you start or restart a task, you can choose whether to start from a saved snapshot state:
| Startup strategy | Description |
|---|---|
| Stateful activation | Restores from an existing valid state (Checkpoint / Savepoint). You can restore from the latest state or select one from the historical snapshot list |
| stateless activate | Starts fresh without any initial state |
Note: Tasks can't be started from a historical version. You can select any Savepoint, but you're responsible for Schema compatibility; Schema changes may cause runtime errors.
4.2.4 Submit and start
The following validations run before startup. The real-time task starts after they pass:
-
Task existence validation
-
Content validity validation
-
Version consistency validation
-
Channel existence validation
-
Snapshot existence validation (for stateful startup)
5. Real-time task running status management
5.1 Task instance management statuses
| Status | Can change to | Available operations | Description |
|---|---|---|---|
| Unpublished | — | Edit, Release, Delete | The real-time task exists only in the development environment |
| Not Started | — | Edit, Release, Delete | Released to production but not yet running |
| Starting | Running, Failed, Stopping, Killing | Edit, Release, Stop, Kill | Submitting the task to the execution engine |
| Running | Stopped, Stopping, Killing, Failed | Edit, Release, Stop, Kill | The task is running normally |
| Failed | Starting (new instance) | Edit, Release | Final state of the instance; a new instance can be started |
| Stopping | Stopped, Killing, Failed | Edit, Release, Kill | Stop operation in progress |
| Stopped | Starting (new instance) | Edit, Release | Final state of the instance; a new instance can be started |
| Killing | Stopped, Failed | Edit, Release | Forcibly terminating the task |
Note: An alert is sent when Stopping or Killing times out. During Stopping, you can stop the task again.
5.2 Differences between Stop and Kill
| Action | Description |
|---|---|
| Stop | Graceful stop. You can choose whether to trigger a Savepoint to save the state. After stopping, you can restore from the Savepoint |
| Kill | Forced termination without saving the state. Use it when Stop doesn't respond or the task is stuck |
- When stopping a task: a Savepoint is triggered automatically (you don't need to define the storage path)
- When starting/restarting a task: you can choose whether to restore from a Checkpoint/Savepoint
6. Production environment detail page
6.1 Instance content
Latest instance information, including
- Script SQL content
- Runtime configuration of the instance
- Table lineage of the instance
- Task parameter values in the instance
6.2 Operation record
Operation records track the full lifecycle events of each instance by Streampark Job Instance ID:
- By default, the event messages of the latest instance are shown, sorted by instance creation time in descending order
- You can expand to view the event messages of all instances
- A timeline controls the display of global information
6.3 Logs
| Category | Content | Description | |
|---|---|---|---|
Run logs | JobManager logs | Task scheduling, Checkpoints, state management | Logs of a single JM instance |
| TaskManager logs | Data processing, network communication | Aggregated logs of multiple TM instances | |
| Exception information | Aggregated ERROR/WARN logs | Entry point for quickly locating problems |
6.4 Snapshot management
| Field | Description |
|---|---|
| Snapshot ID | Checkpoint ID returned by Streampark |
| Status | Current status of the snapshot |
| Snapshot type | Checkpoint: an "auto-save" created automatically by Flink for failure recovery; lightweight, fast, and automated Savepoint: a "manual save point" triggered by the user for planned stops and restarts (code updates, cluster scaling, version upgrades, and so on) |
| Creation method | Manual (user-initiated, Savepoint only) / Automatic (generated automatically based on the Checkpoint configuration) |
| Trigger time | The time the snapshot was triggered |
| Duration | Time taken to create the snapshot |
| Source instance | ID of the instance the snapshot belongs to |
7. Management page metadata
On the real-time task management list page, you can view the following metadata:
| Field | Description |
|---|---|
| Real-time task name | Unique task identifier |
| Remarks | Task description |
| Status | Published / Pending Publishing |
| Running status | Running status of the latest instance of the current version |
| Engine version | For example, Flink 1.16 |
| Owner | Task owner |
| Online time | Time when the task was released from the development environment to production |
| Start time | Start time of the latest instance |
| Last Modified Time | Time of the last modification |
| Last Modifier | User who made the last modification |
| Action | Edit, status management operations, jump to task details |

