如何实现可转换为基类Stream<U>的泛型对象输入流接口?
听起来你想打造一个泛型对象输入流框架——核心是做一个通用接口/轻量代理,让用户能接入自定义流实现(比如Protobuf消息流),还能灵活构建转换管道,把对象流转成字符串流或者其他类型的流对吧?下面我分享一套可落地的设计思路:
1. 定义核心泛型流接口
首先得把最基础的流接口泛型化,既要支持读取原始对象,也要支持链式转换,这样才能构建管道。比如:
public interface ObjectInputStream<T> { // 读取下一个对象,返回null表示流结束 T readNext() throws IOException; // 转换流类型,返回新的泛型流,支持链式调用 <R> ObjectInputStream<R> convert(StreamConverter<T, R> converter); // 可选:批量读取方法,提升效率 List<T> readBatch(int batchSize) throws IOException; }
这里的StreamConverter是专门的转换接口,让用户自定义转换逻辑:
@FunctionalInterface public interface StreamConverter<From, To> { To convert(From input) throws IOException; }
2. 实现轻量级代理基类
为了降低用户的实现成本,你可以提供一个代理基类,用户只需要实现核心的readNext方法,其他通用逻辑(比如转换、批量读取)由基类来处理:
public abstract class AbstractObjectInputStream<T> implements ObjectInputStream<T> { @Override public <R> ObjectInputStream<R> convert(StreamConverter<T, R> converter) { // 返回一个代理流,把当前流的读取结果传给转换器 return new ObjectInputStream<R>() { @Override public R readNext() throws IOException { T original = AbstractObjectInputStream.this.readNext(); return original != null ? converter.convert(original) : null; } @Override public <R2> ObjectInputStream<R2> convert(StreamConverter<R, R2> nextConverter) { return this.convert(nextConverter); } @Override public List<R> readBatch(int batchSize) throws IOException { List<T> originalBatch = AbstractObjectInputStream.this.readBatch(batchSize); return originalBatch.stream() .map(converter::convert) .collect(Collectors.toList()); } }; } @Override public List<T> readBatch(int batchSize) throws IOException { List<T> batch = new ArrayList<>(batchSize); for (int i = 0; i < batchSize; i++) { T obj = readNext(); if (obj == null) break; batch.add(obj); } return batch; } }
3. 用户如何接入自定义流
用户只需要继承这个基类,实现自己的readNext逻辑就行。比如一个Protobuf消息流的实现:
public class ProtobufMessageStream extends AbstractObjectInputStream<MyProtobufMessage> { private final InputStream rawInputStream; public ProtobufMessageStream(InputStream rawInputStream) { this.rawInputStream = rawInputStream; } @Override public MyProtobufMessage readNext() throws IOException { // 这里实现Protobuf的反序列化逻辑 int length = rawInputStream.readInt(); if (length == -1) return null; byte[] buffer = new byte[length]; rawInputStream.readFully(buffer); return MyProtobufMessage.parseFrom(buffer); } }
4. 构建转换管道的示例
用户拿到自己的流后,就可以链式调用convert方法构建管道了,比如把Protobuf流转成字符串流,再转成JSON对象流:
// 1. 创建用户自定义的Protobuf流 ObjectInputStream<MyProtobufMessage> protoStream = new ProtobufMessageStream(new FileInputStream("data.proto")); // 2. 构建转换管道:Protobuf -> 字符串 -> JSON对象 ObjectInputStream<JsonObject> jsonStream = protoStream .convert(protoMsg -> protoMsg.toJsonString()) // 转字符串 .convert(jsonStr -> new JsonParser().parse(jsonStr).getAsJsonObject()); // 转JSON对象 // 3. 读取转换后的流 JsonObject jsonObj; while ((jsonObj = jsonStream.readNext()) != null) { // 处理JSON对象 }
5. 额外的扩展建议
- 可以添加
close方法到接口中,确保流资源能正确释放,基类可以默认实现代理的close逻辑,调用原始流的close。 - 支持异步读取的话,可以在接口中加入
CompletableFuture<T> readNextAsync()方法,适合高并发场景。 - 提供一些内置的转换器,比如
toStringConverter、toByteArrayConverter,方便用户快速使用。
这样设计下来,你的框架既保持了泛型的灵活性,又降低了用户的接入成本,完全满足“自定义转换逻辑+构建转换管道”的需求~
内容的提问来源于stack exchange,提问作者Zelta
相关产品推荐
相关产品推荐

