The pulsar client generics have an exception #19946
|
The code is shown below public class Main {
public static void main(String[] args) throws PulsarClientException {
try (PulsarClient pulsarClient = PulsarClient.builder()
.serviceUrl("pulsar://livk.com:6650")
.build()) {
SchemaDefinition<PulsarMessage<String>> schemaDefinition = SchemaDefinition.<PulsarMessage<String>>builder()
.withPojo(PulsarMessage.class)
.build();
try (Producer<PulsarMessage<String>> producer = pulsarClient.<PulsarMessage<String>>newProducer(Schema.JSON(schemaDefinition))
.topic("livk-topic")
.compressionType(CompressionType.LZ4)
.sendTimeout(0, TimeUnit.SECONDS)
.enableBatching(true)
.batchingMaxPublishDelay(10, TimeUnit.MILLISECONDS)
.batchingMaxMessages(1000)
.maxPendingMessages(1000)
.blockIfQueueFull(true)
.roundRobinRouterBatchingPartitionSwitchFrequency(10)
.batcherBuilder(BatcherBuilder.DEFAULT)
.create()) {
PulsarMessage<String> message = new PulsarMessage<>(LocalDateTime.now(), UUID.randomUUID().toString());
MessageId messageId = producer.newMessage()
.value(message)
.send();
System.out.println(messageId);
}
}
}
@Data
@NoArgsConstructor
@AllArgsConstructor
public static class PulsarMessage<T> {
private LocalDateTime time;
private T data;
}
}Exception in thread "main" org.apache.pulsar.shade.org.apache.avro.AvroTypeException: Unknown type: T
at org.apache.pulsar.shade.org.apache.avro.specific.SpecificData.createSchema(SpecificData.java:403)
at org.apache.pulsar.shade.org.apache.avro.reflect.ReflectData.createSchema(ReflectData.java:780)
at org.apache.pulsar.shade.org.apache.avro.reflect.ReflectData.createFieldSchema(ReflectData.java:873)
at org.apache.pulsar.shade.org.apache.avro.reflect.ReflectData$AllowNull.createFieldSchema(ReflectData.java:92)
at org.apache.pulsar.shade.org.apache.avro.reflect.ReflectData.createSchema(ReflectData.java:736)
at org.apache.pulsar.shade.org.apache.avro.specific.SpecificData$3.computeValue(SpecificData.java:328)
at org.apache.pulsar.shade.org.apache.avro.specific.SpecificData$3.computeValue(SpecificData.java:325)
at java.base/java.lang.ClassValue.getFromHashMap(ClassValue.java:228)
at java.base/java.lang.ClassValue.getFromBackup(ClassValue.java:210)
at java.base/java.lang.ClassValue.get(ClassValue.java:116)
at org.apache.pulsar.shade.org.apache.avro.specific.SpecificData.getSchema(SpecificData.java:339)
at org.apache.pulsar.client.impl.schema.util.SchemaUtil.extractAvroSchema(SchemaUtil.java:95)
at org.apache.pulsar.client.impl.schema.util.SchemaUtil.createAvroSchema(SchemaUtil.java:78)
at org.apache.pulsar.client.impl.schema.util.SchemaUtil.parseSchemaInfo(SchemaUtil.java:51)
at org.apache.pulsar.client.impl.schema.JSONSchema.of(JSONSchema.java:91)
at org.apache.pulsar.client.impl.PulsarClientImplementationBindingImpl.newJSONSchema(PulsarClientImplementationBindingImpl.java:222)
at org.apache.pulsar.client.api.Schema.JSON(Schema.java:362)When I use a class with generics as the message body, the pulsar client cannot parse it jdk: 17 pulsar-client:2.11.0 |
Replies: 2 comments 3 replies
|
Hi @livk-cloud! public static class StringPulsarMessage extends PulsarMessage<String> {
}
SchemaDefinition<StringPulsarMessage> schemaDefinition = SchemaDefinition.<StringPulsarMessage>builder()
.withPojo(StringPulsarMessage.class)
.build();Please let me know if this works for you. |
|
Hi @livk-cloud! Cross refer to:
Quick answer: you may de-parameterize |
Hi @livk-cloud!
Cross refer to:
Quick answer: you may de-parameterize
PulsarMessage<T>toPulsarMessageand inline the type or useObject. This seems an Avro upstream issue.