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
相关产品推荐
相关产品推荐

