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
- CC Flink ≠ Apache Flink. Verify every API, SQL construct, and runtime behavior against:
- Confluent Cloud Flink docs
- CC Flink SQL reference
- Confluent Terraform provider
- Live
confluent flink shellagainst the user's compute pool
- No mocks in verification. Integration claims require real
confluentCLI runs. Unit tests may mock; anything calling itself "end-to-end verification" may not. - Record decisions somewhere durable. Ask the user where they want dialect traps and verification notes tracked (e.g.
docs/flink-dialect-traps.mdin their project) before writing any new file — don't assume adocs/layout. - Secrets never in repo.
CREATE CONNECTIONparameters are Terraform-injected, never hardcoded. Gitignore.tfvars,.tfstate*,*.secret*from day one. - EXPLAIN before CREATE. Always
EXPLAINa query beforestatement create— catches parse/type errors without consuming CFUs. - 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:
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 creatableSET 'execution.checkpointing.*'— CC-managed, not settableCREATE 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 AFTERINSERT INTOGROUP BY TUMBLE(ts, INTERVAL ...)— must use the TVF form:TUMBLE(TABLE t, DESCRIPTOR(ts), ...)LATERAL TABLE(UNNEST(...))— parse error; useCROSS JOIN UNNEST(...)PROCTIME()— not supported; use External Tables/KEY_SEARCH_AGGor 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 '...'withoutUSING JAR— CC UDFs require an uploaded artifact- Savepoints /
STOP WITH SAVEPOINT— not exposed on CC --sql-fileflag — 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:
- EXPLAIN the query in
flink shell— catches syntax and type errors for free. - Write a minimal reproducer.
- 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.
- Run:
confluent flink statement create <name> --sql "$(cat repro.sql)" --compute-pool <id> --database <cluster> --environment <env> --wait - 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'svalue.format(see references/formats-and-serialization.md [blocked];jsonschemaforjson-registry,avroforavro-registry,protobufforproto-registry,stringforraw) - 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 CONNECTIONSclause
See references/cli-reference.md [blocked] for full flag schemas and timing expectations.
Anti-patterns
- Apache Flink docs or Stack Overflow answers tagged
apache-flinkcited as CC authority - LLM memory of "Flink SQL syntax" used without CC verification
- Mocking the
confluentCLI in anything claiming end-to-end verification terraform apply -auto-approveon 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 CONNECTIONor 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
references/dialect-traps.md(this skill) — canonical, consolidated- Per-project
CLAUDE.md— references this skill, adds project-specific context - 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.


