【发布时间】:2020-06-16 09:25:48
【问题描述】:
我正在编写某种带有缓存的中间件 HTTP 代理。工作流程是:
- 客户端向此代理请求资源
- 如果缓存中存在资源,则代理返回它
- 如果未找到资源,则代理获取远程资源并返回给用户。代理在数据加载时将此资源保存到缓存中。
我的接口有用于远程资源的Publisher<ByteBuffer> 流,接受Publisher<ByteBuffer> 保存的缓存,以及接受Publisher<ByteBuffer> 作为响应的客户端连接:
// remote resource
interface Resource {
Publisher<ByteBuffer> fetch();
}
// cache
interface Cache {
Completable save(Publisher<ByteBuffer> data);
}
// clien response connection
interface Connection {
Completable send(Publisher<ByteBuffer> data);
}
我的问题是我需要在向客户端发送响应时延迟保存这个字节缓冲区流以缓存,因此客户端应该负责从远程资源请求ByteByffer块, 不缓存。
我尝试使用Publisher::cache 方法,但这对我来说不是一个好的选择,因为它将所有接收到的数据保存在内存中,这是不可接受的,因为缓存的数据可能只有几 GB 大小。
作为一种解决方法,我创建了Subject,并由来自Resource 的下一个项目填充:
private final Cache cache;
private final Connection out;
Completable proxy(Resource res) {
Subject<ByteBuffer> mirror = PublishSUbject.create();
return Completable.mergeArray(
out.send(res.fetch().doOnNext(mirror::onNext),
cache.save(mirror.toFlowable(BackpressureStrategy.BUFFER))
);
}
是否可以重复使用相同的 Publisher 而无需在内存中缓存项目,并且只有一个订阅者负责从发布者请求项目?
【问题讨论】:
-
我听不懂你说:所以客户端应该负责从远程资源请求ByteByffer块,而不是缓存。!
-
@bubbles 我的意思是背压:
Publisher的Subscription具有request(long n)方法,它正在从Publisher请求下一个n的项目数量。客户端的Connection比Cache慢,所以只有Connection应该负责从远程Publisher的Resource请求下一个n数量的ByteBuffer -
这是什么版本的 RxJava?我在 2.x 上,
Publisher只有一种方法:Publisher.subscribe(Subscriber<? super T>)
标签: java caching rx-java reactive-streams backpressure