Skip to content

[Feature-18691][task-flink] Submit Flink SQL through SQL Gateway - #18693

Open
macdoor wants to merge 1 commit into
apache:devfrom
macdoor:feature/flink-sql-submit
Open

macdoor wants to merge 1 commit into
apache:devfrom
macdoor:feature/flink-sql-submit

Conversation

@macdoor

@macdoor macdoor commented Oct 10, 2026 •

Copy link
Copy Markdown
Contributor

Was this PR generated or assisted by AI?

YES — the SQL Gateway submit path, UI fields, and unit tests were assisted by an AI agent (Hermes). All changes were human-reviewed and verified by build + unit tests plus a smoke run on an internal cluster.

Purpose of the pull request

Add an optional SQL Gateway submit mode to the existing FLINK task, so batch SQL can be sent to a remote Flink SQL Gateway through its JDBC driver (flink-sql-jdbc-driver), enabling DolphinScheduler to schedule batch SQL workloads without a local Flink client installation.

Today the FLINK task runs SQL through ${FLINK_HOME}/bin/sql-client.sh, which requires a full Flink distribution on every worker host. Deployments that run a dedicated Flink SQL Gateway cluster instead can submit scripts straight to the Gateway over JDBC.

Compared with spawning sql-client.sh on the worker, submitting through the SQL Gateway costs less memory on the worker host (no local Flink CLI/mini-cluster JVM per task), produces clearer errors (statement-level JDBC exceptions instead of parsing shell/YARN logs), and is faster (no CLI process startup and script materialization per statement).

The new mode is opt-in via a sqlSubmitType parameter: absent or CLIENT keeps the existing sql-client.sh path unchanged; SQL_GATEWAY submits over JDBC.

Brief change log

  • FlinkParameters.sqlSubmitType: absent / CLIENT keeps ${FLINK_HOME}/bin/sql-client.sh; SQL_GATEWAY submits over JDBC and cancels via statement.cancel()
  • Fork only at FlinkTask.handle() / cancelApplication(). getScript() and FlinkArgsUtils are unchanged, so existing CLIENT behavior and FLINK_STREAM are untouched
  • Carries the same capability the SQL task gained in [Improvement-18019][task-sql] Support SQL from resource file and parameter placeholders #18020 — SQL from a resource-center file and parameter placeholders — into the FLINK task: script sources reuse org.apache.dolphinscheduler.plugin.task.api.enums.SqlSourceType (SCRIPT / FILE, introduced by [Improvement-18019][task-sql] Support SQL from resource file and parameter placeholders #18020) for both the init script and the main script, and the JDBC URL and script content resolve parameter placeholders (e.g. ${var}) before submission; statement splitting and leading-comment stripping are package-private static methods covered by plain JUnit tests (no Mockito)
  • Statement separator is configurable (default ;); optional maxPrintRows caps query-result row logging; flinkJdbcUrl supports parameter placeholders; optional jdbcProperties map filters null keys/values before connect
  • The JDBC driver is shaded only into dolphinscheduler-task-flink; dolphinscheduler-task-flink-stream excludes it so its shade jar stays free of the driver
  • UI: the SQL Gateway fields show only when program type is SQL and submit type is SQL Gateway; use-resources field helper gains optional field / limit parameters to support the init-script resource field; locales (en/zh)

Verify this pull request

This change added tests and can be verified as follows:

  • mvn test -pl dolphinscheduler-task-plugin/dolphinscheduler-task-flink -Dtest=FlinkParametersTest,FlinkSqlGatewayExecutorTest,FlinkArgsUtilsTest,FlinkTaskTest -Dspotless.skip=true -Djacoco.skip=true
  • mvn -pl dolphinscheduler-task-plugin/dolphinscheduler-task-flink,dolphinscheduler-task-plugin/dolphinscheduler-task-flink-stream spotless:check
  • A FLINK task with sqlSubmitType=SQL_GATEWAY was executed on an internal cluster (worker-only smoke): the task instance ran FlinkSqlGatewayExecutor, connected to jdbc:flink://...:8083, executed init + main scripts from resource-center files, and cancelled via statement.cancel()

Keep sql-client.sh when sqlSubmitType is absent or CLIENT. SQL_GATEWAY
submits over JDBC and cancels the statement.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

backend test UI ui and front end related

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant