Skip to main content

Real-time tasks

Last updated 10/03/2026

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.

FieldDescriptionRemarks
NameUnique within the workspace, up to 48 characters; supports Chinese characters, English letters, digits, and underscoresRequired
Flink versionThe Flink engine version bound to the real-time task, selected from the drop-down listCan't be changed after creation
OwnerThe owner of the real-time task. The creator is the owner by defaultRequired
RemarksTask description, up to 200 charactersOptional

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 typeDescription
Workspace ParamsSupports 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):

  1. 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
  1. 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:

  1. SQL syntax validation

  2. Table permission validation (read and write permissions on the source and target tables)

  3. Parameter parsing and recording

Debug/run validation (in addition to save validation):

  1. Parameter completeness validation — all ${XX} parameters must have values
  2. Workspace parameter existence validation
  3. 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":

ItemDescription
EnvironmentSaved only in the Dev environment
StatusPending Publishing
VersionThe development environment shows the latest version; none in production
Available operationsEdit, Release, Delete

Saving a task that's "Published":

ItemDescription
EnvironmentThe new version is saved only in the Dev environment and must be released to production again
StatusPublished (with changes)
VersionThe development environment shows the latest version; production still has the previously released version
Available operationsEdit, 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:

  1. Unsaved tasks can't be released (the release button is grayed out)
  2. Production environment table permission validation
  3. Parameter value completeness validation
  4. 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:

ModuleDescription
Real-time task basic informationTask name, owner, engine version, change status, creation/modification time, description
Operation settingExecution channel selection (resource configuration), runtime configuration parameters
Instance contentSQL content, runtime configuration, table lineage, and task parameter values of the latest instance
Operation recordEvent messages for the full lifecycle of instances
LogIncludes run logs and exception logs
SnapshotTotal 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 itemsDescription

Channel (resource pool) selection

Channel name: the system automatically filters the available channels based on the task's Flink engine version

Only "Enabled" channels can be selected

Real-time tasks can only run on channels (resource pools) with Run mode = Application

Calculate mode
  • Resource modes: Compute Optimized / General Purpose.
  • Compute Optimized (1:2): 2G memory per core, stronger CPU performance, suitable for intensive computing; General Purpose (1:4): 4G memory per core, balanced CPU and memory, suitable for most workloads
  • The available resource modes depend on the specs bound to the selected resource pool
Parallelism

Resource specs are defined at the channel level. When an instance is submitted for execution, you choose the parallelism.

Parallelism describes how Flink splits a task for parallel execution. The number of slots is the granularity at which a single TaskManager's resources are divided.

Runtime parameter configuration

CategoryParameterDescriptionDefault value
Checkpoint settingsexecution.checkpointing.intervalCheckpoint interval180 seconds
execution.checkpointing.timeoutCheckpoint timeout3 minutes
execution.checkpointing.min-pauseMinimum pause between checkpoints60 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 strategyrestart-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:

  • Failure Rate (restart based on failure rate): set the detection interval, the maximum number of failures, and the interval between restarts
  • Fixed Delay (restart at a fixed interval): set the number of restart attempts and the interval between restarts
  • No Restarts (no restart): the job doesn't restart automatically after it fails
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 strategyDescription
Stateful activationRestores from an existing valid state (Checkpoint / Savepoint). You can restore from the latest state or select one from the historical snapshot list
stateless activateStarts 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:

  1. Task existence validation

  2. Content validity validation

  3. Version consistency validation

  4. Channel existence validation

  5. Snapshot existence validation (for stateful startup)

5. Real-time task running status management​

5.1 Task instance management statuses​

StatusCan change toAvailable operationsDescription
Unpublished—Edit, Release, DeleteThe real-time task exists only in the development environment
Not Started—Edit, Release, DeleteReleased to production but not yet running
StartingRunning, Failed, Stopping, KillingEdit, Release, Stop, KillSubmitting the task to the execution engine
RunningStopped, Stopping, Killing, FailedEdit, Release, Stop, KillThe task is running normally
FailedStarting (new instance)Edit, ReleaseFinal state of the instance; a new instance can be started
StoppingStopped, Killing, FailedEdit, Release, KillStop operation in progress
StoppedStarting (new instance)Edit, ReleaseFinal state of the instance; a new instance can be started
KillingStopped, FailedEdit, ReleaseForcibly 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​

ActionDescription
StopGraceful stop. You can choose whether to trigger a Savepoint to save the state. After stopping, you can restore from the Savepoint
KillForced 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​

CategoryContentDescription

Run logs

JobManager logsTask scheduling, Checkpoints, state managementLogs of a single JM instance
TaskManager logsData processing, network communicationAggregated logs of multiple TM instances
Exception informationAggregated ERROR/WARN logsEntry point for quickly locating problems

6.4 Snapshot management​

FieldDescription
Snapshot IDCheckpoint ID returned by Streampark
StatusCurrent 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 methodManual (user-initiated, Savepoint only) / Automatic (generated automatically based on the Checkpoint configuration)
Trigger timeThe time the snapshot was triggered
DurationTime taken to create the snapshot
Source instanceID 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:

FieldDescription
Real-time task nameUnique task identifier
RemarksTask description
StatusPublished / Pending Publishing
Running statusRunning status of the latest instance of the current version
Engine versionFor example, Flink 1.16
OwnerTask owner
Online timeTime when the task was released from the development environment to production
Start timeStart time of the latest instance
Last Modified TimeTime of the last modification
Last ModifierUser who made the last modification
ActionEdit, status management operations, jump to task details
Was this page helpful?