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

Apache Beam Pulsar IO 2.69.0:PulsarMessage私有getter访问方法咨询

Apache Beam Pulsar IO 2.69.0:PulsarMessage getter方法访问权限问题

问题描述

在Apache Beam Pulsar IO连接器2.69.0版本中,PulsarMessage类的所有getter方法(如getTopic()、getValue())都是包私有访问权限,导致外部包无法直接调用这些方法获取消息的topic名称和payload,报错信息如下:

'getValue()' is not public in 'org.apache.beam.sdk.io.pulsar.PulsarMessage'. Cannot be accessed from outside package

PulsarMessage抽象类代码

@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<String, String> getProperties();

  @SuppressWarnings("mutable")
  abstract byte[] getMessageId();

  public static PulsarMessage create(
      @Nullable String topicName,
      long publishTimestamp,
      @Nullable String key,
      byte[] value,
      @Nullable Map<String, String> properties,
      byte[] messageId) {
    return new AutoValue_PulsarMessage(
        topicName, publishTimestamp, key, value, properties, messageId);
  }

  public static PulsarMessage create(Message<byte[]> message) {
    return create(
        message.getTopicName(),
        message.getPublishTime(),
        message.getKey(),
        message.getValue(),
        message.getProperties(),
        message.getMessageId().toByteArray());
  }
}

AutoValue生成的实现类代码

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<String, String> 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<String, String> getProperties() {
    return this.properties;
  }

  byte[] getMessageId() {
    return this.messageId;
  }
}

触发错误的测试代码

static class LogMessageFn extends DoFn<PulsarMessage, Void> {
    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()); // 触发权限错误
      } catch (Exception e){
        System.out.println("Exception");
      }
    }
  }

问题解答

这是有意设计吗?

不是,这明显是一个bug。PulsarMessage作为对外暴露的消息载体类,其核心属性的getter方法理应是public权限。问题根源在于抽象类中未给getter方法添加public修饰符,导致AutoValue生成的实现类也沿用了包私有权限,外部包无法访问。

如何访问消息的topic和payload?

有两种可行方案:

  1. 升级到修复后的版本
    Beam团队在后续版本(如2.70.0及以后)已修复此问题,将PulsarMessage的所有getter方法改为public权限,直接升级依赖版本即可解决。

  2. 临时兼容方案(针对2.69.0版本)

    • 反射调用getter:通过Java反射绕过访问权限限制,示例代码如下:
      @ProcessElement
      public void processElement(@Element PulsarMessage message) {
        try {
          // 获取payload
          Method getValueMethod = PulsarMessage.class.getDeclaredMethod("getValue");
          getValueMethod.setAccessible(true);
          byte[] payload = (byte[]) getValueMethod.invoke(message);
          System.out.println("message value: " + new String(payload));
      
          // 获取topic
          Method getTopicMethod = PulsarMessage.class.getDeclaredMethod("getTopic");
          getTopicMethod.setAccessible(true);
          String topic = (String) getTopicMethod.invoke(message);
          System.out.println("message topic: " + topic);
        } catch (Exception e) {
          e.printStackTrace();
        }
      }
      
    • 直接处理原始Pulsar Message:读取消息时使用PulsarIO.readMessagesWithMetadata()获取原始Message<byte[]>对象,直接提取属性:
      PCollection<Message<byte[]>> messages = pipeline.apply(
          PulsarIO.readMessagesWithMetadata()
              .serviceUrl("pulsar://localhost:6650")
              .topic("your-topic")
              .subscriptionName("your-subscription")
      );
      
      messages.apply(ParDo.of(new DoFn<Message<byte[]>, Void>() {
        @ProcessElement
        public void processElement(@Element Message<byte[]> message) {
          String topic = message.getTopicName();
          byte[] payload = message.getValue();
          System.out.println("topic: " + topic + ", payload: " + new String(payload));
        }
      }));
      

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.11 22:13:18