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?
有两种可行方案:
升级到修复后的版本
Beam团队在后续版本(如2.70.0及以后)已修复此问题,将PulsarMessage的所有getter方法改为public权限,直接升级依赖版本即可解决。临时兼容方案(针对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)); } }));
- 反射调用getter:通过Java反射绕过访问权限限制,示例代码如下:
内容的提问来源于stack exchange,提问作者Vaibhav Chandra
相关产品推荐
相关产品推荐

