Flink Udf

作者 confluentinc914d95eff7ff无许可证58 个星标收录于 2026年10月8日更新于 2026年10月8日仓库2天前更新

Build and deploy Apache Flink user-defined functions (UDFs) in Java for stream processing over Kafka. Use this skill when users want to create scalar UDFs, user-defined table functions (UDTFs), or process table functions (PTFs) in Java, deploy them to Confluent Cloud or local Docker environments, and invoke them from Flink SQL or the Table API. Trigger on: Flink UDF, custom Flink function, process table function, PTF, UDTF, Flink user defined, extend Flink SQL, stateful stream processing with Flink. Do NOT trigger for: Kafka Streams UDFs (use kafka-streams-programming skill), general Flink job development without custom functions, CDC streaming data piplines that include Flink (prefer the confluent-cloud-cdc-tableflow skill), Flink connector setup, or Kafka producer/consumer code.

AI 生成的概览

指导用 Java 构建并部署 Apache Flink 用户自定义函数(UDF、UDTF、PTF)到 Confluent Cloud 或本地 Docker。

功能
该技能引导用户为 Apache Flink 创建自定义 Java 函数:标量 UDF、用户自定义表函数(UDTF)和过程表函数(PTF)。它会询问部署目标、使用现有还是新建基础设施以及调用方式,然后指向 Confluent Cloud 或本地 Docker 的参考指南。内容涵盖生成样板代码、实现业务逻辑、打包 JAR、部署制品、在 Flink 中注册函数并用示例数据测试。它还要求在运行会修改资源的命令前先给出部署计划并获得用户明确批准。
适用场景
当用户想用 Java 编写或部署自定义 Flink 函数(如标量 UDF、UDTF 或 PTF),并从 Flink SQL 或 Table API 调用时使用。适用于 Confluent Cloud 或本地 Docker 部署,包括需要先搭建 Kafka 与 Flink 基础设施的情况。不适用于 Kafka Streams UDF、不含自定义函数的通用 Flink 作业开发、CDC 流水线、连接器配置或 Kafka 生产者/消费者代码。
运行要求
需要一个能读取随附参考 Markdown 文件的智能体;不包含脚本。部署路径假定具备用于打包 JAR 的 Java 构建工具、带 confluent CLI 的 Confluent Cloud 环境与计算池,或包含 Kafka 和 Flink 容器的本地 Docker 环境,并需要访问这些服务的网络。

Flink User-Defined Functions (UDFs)

Build and deploy custom functions in Java for Apache Flink to extend SQL and Table API capabilities with custom logic.

Function Types

Before proceeding, identify which type of function the user needs:

  • Scalar UDF: Maps input values to a single output value (e.g., custom hash, string manipulation, calculations)
  • User-Defined Table Function (UDTF): Maps input to multiple output rows (e.g., split strings, explode arrays)
  • Process Table Function (PTF): Advanced stateful processing with N-to-M semantics, managed state, and timers (e.g., windowing, deduplication, state machines)

Gather Requirements

Ask the user these questions to determine the implementation path (if not already clear from context):

  1. Deployment target: Confluent Cloud or local Docker?
  2. Infrastructure: Deploy new infrastructure (Kafka + Flink) or use existing?
  3. Invocation method: Flink SQL or Table API?

Route to Implementation Guide

Based on the answers above, read the appropriate reference file:

Confluent Cloud Deployment

  • Scalar UDF or UDTF → Read references/udf-udtf-java-confluent-cloud.md
  • Process Table Function (PTF) → Read references/ptf-java-confluent-cloud.md

If infrastructure setup is needed, also read references/confluent-cloud-setup.md first.

Local Docker Deployment

  • Scalar UDF or UDTF → Read references/udf-udtf-java-local.md
  • Process Table Function (PTF) → Read references/ptf-java-local.md

If infrastructure setup is needed, also read references/local-docker-setup.md first.

Implementation Workflow

After reading the appropriate reference:

  1. Set up infrastructure (if needed)
  2. Generate boilerplate code for the function
  3. Implement the business logic
  4. Build and package the JAR
  5. Confirm the deployment plan with the user. Before any resource-modifying call (confluent flink artifact create, docker cp into a running container, CREATE FUNCTION, etc.), present the plan and wait for explicit approval. Show:
    • Artifact name and JAR path
    • Function name to register
    • Target environment (Confluent Cloud env + compute pool ID, or local Docker container name)
    • The exact commands and SQL that will run Do not proceed to steps 6–7 until the user confirms.
  6. Deploy the artifact
  7. Register the function in Flink
  8. Test the function with sample data
  9. Provide usage examples (SQL or Table API)

Keep code scaffolding concise and focused on the user's specific requirements. Avoid over-engineering.

来源与署名

来源:confluentinc/agent-skills位于skills/flink-udf提交914d95e

许可证: 无许可证

内容归原作者所有。SourceWeft 从公开仓库中收录这些内容。

举报或申请下架