Confluent Cloud Flink Sql

作者 confluentinc914d95eff7ff無授權條款58 個星標收錄於 2026年10月8日更新於 2026年10月8日儲存庫2 天前更新

Write and debug Flink SQL that runs on Confluent Cloud, enforcing the CC-vs-Apache-Flink (OSS) dialect boundary. Use when the working directory is a Confluent Cloud Flink workspace, when a Flink SQL statement needs checking before it runs on a CC compute pool, or when the user mentions CC Flink, Confluent Cloud Flink SQL, the `confluent flink` CLI, a CFU compute pool, `CREATE CONNECTION`, or asks to check or debug Flink SQL whose runtime is Confluent Cloud. Also trigger when a Flink SQL question is posed and nothing establishes an Apache Flink OSS runtime. Do NOT trigger for: building or deploying Flink UDFs in Java (UDF/UDTF/PTF — use flink-udf); a full CDC pipeline from a database through Flink into Tableflow/Iceberg/Delta Lake (use confluent-cloud-cdc-tableflow); Kafka Streams topology work (use kafka-streams-programming); or Flink SQL confirmed to run on Apache Flink OSS, not Confluent Cloud.

AI 產生的概覽

指導撰寫與偵錯在 Confluent Cloud 上執行的 Flink SQL,並強制區分 CC 與 Apache Flink 方言。

功能
此技能提供撰寫、檢查與偵錯在 Confluent Cloud 上執行之 Flink SQL 的說明與參考資料。內容涵蓋方言陷阱、CLI 參數結構、經驗證的 SQL 模式、支援的格式、保留字,以及 CC 特有的錯誤修正,並規定使用 confluent CLI 的「先 EXPLAIN 再建立」驗證流程。產出為修正後的 SQL、驗證指令與記錄的輸出,而非檔案或指令碼。
適用情境
適用於在 Confluent Cloud Flink 工作區中作業、需要在 CC 運算池上執行前檢查 Flink SQL 語句,或使用者提到 CC Flink、confluent flink CLI、CFU 運算池或 CREATE CONNECTION 的情況。不適用於 Apache Flink OSS 執行環境、Java UDF 開發、Kafka Streams 拓撲,或寫入 Tableflow 的完整 CDC 管線。
執行需求
需要已驗證工作階段的 confluent CLI,以及使用中的 Confluent Cloud 運算池;執行語句會消耗 CFU。Terraform 為選用項目,僅在透過 Confluent Terraform provider 管理 CREATE CONNECTION 認證時需要。隱含需要連線至 Confluent Cloud 及其文件的網路存取。不附指令碼,僅有說明與參考文件。

Confluent Cloud Flink SQL

Enforce the CC-Flink-vs-OSS-Flink dialect boundary and the CLI-driven verification loop for any Confluent Cloud Flink SQL work. Apache Flink OSS training data is a trap — CC rejects or silently mishandles a long list of otherwise-valid Flink SQL constructs.

Scope note: this skill's reference material is built around the OSS-vs-CC dialect boundary — traps in constructs that exist in both dialects but behave differently. It does not yet catalog CC-only DDL that has no OSS counterpart (e.g. CREATE MATERIALIZED TABLE, CREATE MODEL/AI_COMPLETE, CREATE AGENT, USE CATALOG). For those, verify directly against the CC Flink SQL reference rather than expecting a trap entry here.

Non-negotiables

  1. CC Flink ≠ Apache Flink. Verify every API, SQL construct, and runtime behavior against:
  2. No mocks in verification. Integration claims require real confluent CLI runs. Unit tests may mock; anything calling itself "end-to-end verification" may not.
  3. Record decisions somewhere durable. Ask the user where they want dialect traps and verification notes tracked (e.g. docs/flink-dialect-traps.md in their project) before writing any new file — don't assume a docs/ layout.
  4. Secrets never in repo. CREATE CONNECTION parameters are Terraform-injected, never hardcoded. Gitignore .tfvars, .tfstate*, *.secret* from day one.
  5. EXPLAIN before CREATE. Always EXPLAIN a query before statement create — catches parse/type errors without consuming CFUs.
  6. Don't invent identifiers. Use <placeholder> for any topic, table, statement, or resource name you haven't verified.

Reference files

Load these on demand when the topic matches — do not read them all upfront:

FileWhen to load
references/dialect-traps.md [blocked]Before writing ANY Flink SQL — 22 CC-vs-OSS traps, single source of truth
references/cli-reference.md [blocked]Before running confluent CLI — flag schemas, carry-over recipe, timing, token expiry
references/sql-patterns-cc.md [blocked]When writing SQL — CC-validated patterns: windows, joins, dedup, MATCH_RECOGNIZE, JSON, External Tables
references/formats-and-serialization.md [blocked]When configuring table formats — 7 supported formats, id-encoding, consume flags
references/troubleshooting-cc.md [blocked]When debugging errors — CC-specific error to cause to fix
references/reserved-words.md [blocked]When hitting parse errors — must-backquote identifiers

Red flags

Stop and consult references/dialect-traps.md if you catch yourself writing any of these:

  • DataStream API (Java/Scala) — not supported on CC. Table API + PTF are GA in Java; Python Table API is Open Preview with no PTF yet
  • CREATE CATALOG ... — catalog = CC environment, not creatable
  • SET 'execution.checkpointing.*' — CC-managed, not settable
  • CREATE TABLE ... WITH ('connector' = 'kafka', ...) — tables auto-map from topics
  • 'value.format' = 'json' — must be 'json-registry' (or another SR-backed format)
  • WITH cte AS (...) INSERT INTO ... — CC requires the CTE AFTER INSERT INTO
  • GROUP BY TUMBLE(ts, INTERVAL ...) — must use the TVF form: TUMBLE(TABLE t, DESCRIPTOR(ts), ...)
  • LATERAL TABLE(UNNEST(...)) — parse error; use CROSS JOIN UNNEST(...)
  • PROCTIME() — not supported; use External Tables/KEY_SEARCH_AGG or an event-time temporal join (avoid a regular join against an upsert-kafka topic — it retains the whole table in state)
  • CREATE FUNCTION f AS '...' without USING JAR — CC UDFs require an uploaded artifact
  • Savepoints / STOP WITH SAVEPOINT — not exposed on CC
  • --sql-file flag — doesn't exist; use --sql "$(cat file.sql)"
  • DROP TABLE — deletes the physical Kafka topic and its data on CC, not just metadata; confirm before running

Verification loop

Canonical validation loop for any CC Flink SQL claim:

  1. EXPLAIN the query in flink shell — catches syntax and type errors for free.
  2. Write a minimal reproducer.
  3. Present the plan and wait for explicit user confirmation before running anything that creates or modifies a real resource. State: the statement name and SQL, the compute pool/database/environment it targets, and any side effects (DDL creates a Kafka topic and Schema Registry subject; every run consumes CFUs). Do not proceed to step 3 without a go-ahead.
  4. Run: confluent flink statement create <name> --sql "$(cat repro.sql)" --compute-pool <id> --database <cluster> --environment <env> --wait
  5. Observe. Consume downstream: confluent kafka topic consume <topic> --cluster <id> --from-beginning --value-format <matching-format> 2>/dev/null | grep -v '^%' — match <matching-format> to the sink's value.format (see references/formats-and-serialization.md [blocked]; jsonschema for json-registry, avro for avro-registry, protobuf for proto-registry, string for raw)
  6. Record the command + output for later reference.

Escalation-required states (no silent workarounds):

  • Statement PENDING > 60s → confluent flink statement exception list <name> --cloud <provider> --region <region>
  • UDF deploy "jar not found" → confluent flink artifact list --cloud <provider> --region <region>
  • Schema mismatch → DESCRIBE <table>, diff against the producer schema
  • Egress denied → check CREATE CONNECTION + USING CONNECTIONS clause

See references/cli-reference.md [blocked] for full flag schemas and timing expectations.

Anti-patterns

  • Apache Flink docs or Stack Overflow answers tagged apache-flink cited as CC authority
  • LLM memory of "Flink SQL syntax" used without CC verification
  • Mocking the confluent CLI in anything claiming end-to-end verification
  • terraform apply -auto-approve on the first run of a root module
  • Committing .tfvars, .tfstate*, .terraform/, *.secret*
  • Swallowing Flink statement exceptions — fail loud; read statement exception list
  • Hardcoded secrets in CREATE CONNECTION or UDF source

Tutorials

developer.confluent.io/tutorials/#flink mixes OSS and CC tutorials.

Filter rule: Only use tutorials that list "Confluent Cloud" in prerequisites or use confluent flink shell. Apply the dialect trap table to any SQL copied from a tutorial — many target OSS Flink or Kafka Streams and are not CC-compatible as-is.

Source-of-truth hierarchy

  1. references/dialect-traps.md (this skill) — canonical, consolidated
  2. Per-project CLAUDE.md — references this skill, adds project-specific context
  3. A per-project trap log, if the user wants one kept — ask where before creating it

When a new trap is discovered during a session, tell the user so they can decide whether to record it in their project's own notes. Do not edit this skill's own installed files (references/dialect-traps.md or elsewhere) — propose the change and let the user (or a separate PR to this skill's repo) apply it.

References

來源與署名

來源:confluentinc/agent-skills位於skills/confluent-cloud-flink-sql提交914d95e

授權條款: 無授權條款

內容歸原作者所有。SourceWeft 從公開儲存庫中收錄這些內容。

檢舉或申請下架