You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

KSQLDB自定义带Struct类型的UDAF遇AnnotationParser异常

自定义KSQLDB UDAF添加SchemaDescriptor后触发空指针异常

问题描述

尝试为KSQLDB编写自定义UDAF,参考官方示例实现。结构体Schema需通过SchemaDescriptor提交给KSQLDB,使其类型系统能够识别。构建shadowJar并部署到KSQLDB的extension_dir目录时:

  • 若不添加param、aggregate、return类型的Schema,KSQLDB可正常启动,且自定义函数能正常显示;
  • 一旦给@UdafFactory注解添加SchemaDescriptor,启动时就会抛出Java空指针异常及AnnotationParser异常,无法解析这个与官方示例一致的Schema。

Schema定义

public static final Schema PARAM_SCHEMA = SchemaBuilder.struct().optional()
    .field("C", Schema.OPTIONAL_INT64_SCHEMA)
    .build();

public static final String PARAM_SCHEMA_DESCRIPTOR = "STRUCT<" +
    "C BIGINT" +
    ">";

public static final Schema AGGREGATE_SCHEMA = SchemaBuilder.struct().optional()
    .field("MIN", Schema.OPTIONAL_INT64_SCHEMA)
    .field("MAX", Schema.OPTIONAL_INT64_SCHEMA)
    .field("COUNT", Schema.OPTIONAL_INT64_SCHEMA)
    .build();

public static final String AGGREGATE_SCHEMA_DESCRIPTOR = "STRUCT<" +
    "MIN BIGINT," +
    "MAX BIGINT," +
    "COUNT BIGINT" +
    ">";

public static final Schema RETURN_SCHEMA = SchemaBuilder.struct().optional()
    .field("MIN", Schema.OPTIONAL_INT64_SCHEMA)
    .field("MAX", Schema.OPTIONAL_INT64_SCHEMA)
    .field("COUNT", Schema.OPTIONAL_INT64_SCHEMA)
    .field("DIFFERENTIAL", Schema.OPTIONAL_INT64_SCHEMA)
    .build();

public static final String RETURN_SCHEMA_DESCRIPTOR = "STRUCT<" +
    "MIN BIGINT," +
    "MAX BIGINT," +
    "COUNT BIGINT," +
    "DIFFERENTIAL BIGINT" +
    ">";

错误日志

java.lang.NullPointerException
    at java.base/sun.reflect.annotation.AnnotationParser.parseArray(AnnotationParser.java:533)
   at java.base/sun.reflect.annotation.AnnotationParser.parseMemberValue(AnnotationParser.java:356)
    at java.base/sun.reflect.annotation.AnnotationParser.parseAnnotation2(AnnotationParser.java:287)
    at java.base/sun.reflect.annotation.AnnotationParser.parseAnnotations2(AnnotationParser.java:121)
    at java.base/sun.reflect.annotation.AnnotationParser.parseAnnotations(AnnotationParser.java:73)
    at java.base/java.lang.reflect.Executable.declaredAnnotations(Executable.java:604)
     at java.base/java.lang.reflect.Executable.declaredAnnotations(Executable.java:602)
    at java.base/java.lang.reflect.Executable.getAnnotation(Executable.java:572)
     at java.base/java.lang.reflect.Method.getAnnotation(Method.java:695)
     at io.confluent.ksql.function.UdafLoader.loadUdafFromClass(UdafLoader.java:59)
    at io.confluent.ksql.function.UserFunctionLoader.loadFunctions(UserFunctionLoader.java:124)
    at io.confluent.ksql.function.UserFunctionLoader.lambda$load$2(UserFunctionLoader.java:97)
     at java.base/java.util.stream.ForEachOps$ForEachOp$OfRef.accept(ForEachOps.java:183)
    at java.base/java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:195)
     at java.base/java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:195)
    at java.base/java.util.stream.ReferencePipeline$2$1.accept(ReferencePipeline.java:177)
    at java.base/java.util.Iterator.forEachRemaining(Iterator.java:133)
    at java.base/java.util.Spliterators$IteratorSpliterator.forEachRemaining(Spliterators.java:1801)
     at java.base/java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:484)
     at java.base/java.util.stream.AbstractPipeline.wrapAndCopyInto(AbstractPipeline.java:474)
     at java.base/java.util.stream.ForEachOps$ForEachOp.evaluateSequential(ForEachOps.java:150)
    at java.base/java.util.stream.ForEachOps$ForEachOp$OfRef.evaluateSequential(ForEachOps.java:173)
    at java.base/java.util.stream.AbstractPipeline.evaluate(AbstractPipeline.java:234)
     at java.base/java.util.stream.ReferencePipeline.forEach(ReferencePipeline.java:497)
     at io.confluent.ksql.function.UserFunctionLoader.load(UserFunctionLoader.java:96)
     at io.confluent.ksql.rest.server.KsqlServerMain.loadFunctions(KsqlServerMain.java:133)
     at io.confluent.ksql.rest.server.KsqlServerMain.main(KsqlServerMain.java:80)

Gradle依赖

dependencies {
implementation "io.confluent.ksql:ksqldb-udf:7.3.0"
implementation "org.apache.kafka:kafka_2.13:2.5.0"
implementation "org.apache.kafka:connect-api:3.3.1"}

内容的提问来源于stack exchange,提问作者Hansanho

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.05 15:00:46