I noticed that the class PulsarMessage has private getters in version 2.69.0. Shouldn't they be public inorder to access the topics names and/or payload of the message. Artifact link : https://mvnrepository.com/artifact/org.apache.beam/beam-sdks-java-io-pulsar
PulsarMessage class:
@DefaultSchema(AutoValueSchema.class)
@AutoValue
public abstract class PulsarMessage {
abstract @Nullable String getTopic();
abstract long getPublishTimestamp();
abstract @Nullable String getKey();
@SuppressWarnings("mutable")
abstract byte[] getValue();
abstract @Nullable Map getProperties();
@SuppressWarnings("mutable")
abstract byte[] getMessageId();
public static PulsarMessage create(
@Nullable String topicName,
long publishTimestamp,
@Nullable String key,
byte[] value,
@Nullable Map properties,
byte[] messageId) {
return new AutoValue_PulsarMessage(
topicName, publishTimestamp, key, value, properties, messageId);
}
public static PulsarMessage create(Message message) {
return create(
message.getTopicName(),
message.getPublishTime(),
message.getKey(),
message.getValue(),
message.getProperties(),
message.getMessageId().toByteArray());
}
}
Error : 'getValue()' is not public in 'org.apache.beam.sdk.io.pulsar.PulsarMessage'. Cannot be accessed from outside package
The concrete implementation provided in the jar
final class AutoValue_PulsarMessage extends PulsarMessage {
private final @Nullable String topic;
private final long publishTimestamp;
private final @Nullable String key;
private final byte[] value;
private final @Nullable Map properties;
private final byte[] messageId;
@Nullable String getTopic() {
return this.topic;
}
long getPublishTimestamp() {
return this.publishTimestamp;
}
@Nullable String getKey() {
return this.key;
}
byte[] getValue() {
return this.value;
}
@Nullable Map getProperties() {
return this.properties;
}
byte[] getMessageId() {
return this.messageId;
}
Is this intended i.e., is there another way to access the topics and payloads?
Below is what I am trying to achieve and facing error
static class LogMessageFn extends DoFn {
private static final long serialVersionUID = 1L;
@ProcessElement
public void processElement(@Element PulsarMessage message) {
try{
System.out.println("message value : " + message.toString());
// System.out.println("message value : " + message.getValue()); throws error
} catch (Exception e){
System.out.println("Exception");
}
}
}