【问题标题】:Shortest code to cache Rxjs http request while not complete?未完成时缓存 Rxjs http 请求的最短代码?
【发布时间】:2016-10-21 14:43:24
【问题描述】:

我正在尝试创建一个满足以下要求的可观察流:

  1. 在订阅时从存储中加载数据
  2. 如果数据尚未过期,则返回存储值的可观察对象
  3. 如果数据过期,则返回一个 HTTP 请求 observable,它使用刷新令牌获取新值 存储它
    • 如果在请求完成之前再次到达此代码,则返回相同的可观察请求
    • 如果在上一个请求完成后或使用不同的刷新令牌到达此代码,则启动新请求

我知道关于如何执行第 (3) 步有很多不同的答案,但是当我尝试一起执行这些步骤时,我正在寻找有关我提出的解决方案是否是它可以是最简洁的(我对此表示怀疑!)。

这是一个演示我当前方法的示例:

var cachedRequestToken;
var cachedRequest;

function getOrUpdateValue() {
    return loadFromStorage()
      .flatMap(data => {
        // data doesn't exist, shortcut out
        if (!data || !data.refreshtoken) 
            return Rx.Observable.empty();

        // data still valid, return the existing value
        if (data.expires > new Date().getTime())
            return Rx.Observable.return(data.value);

        // if the refresh token is different or the previous request is 
        // complete, start a new request, otherwise return the cached request
        if (!cachedRequest || cachedRequestToken !== data.refreshtoken) { 
            cachedRequestToken = data.refreshtoken;

            var pretendHttpBody = {
                value: Math.random(),
                refreshToken: Math.random(),
                expires: new Date().getTime() + (10 * 60 * 1000) // set by server, expires in ten minutes
            };

            cachedRequest = Rx.Observable.create(ob => { 
                // this would really be a http request that exchanges
                // the one use refreshtoken for new data, then saves it
                // to storage for later use before passing on the value

                window.setTimeout(() => { // emulate slow response   
                    saveToStorage(pretendHttpBody);
                    ob.next(pretendHttpBody.value);                 
                    ob.completed(); 
                    cachedRequest = null; // clear the request now we're complete
                }, 2500);
          });
        }

        return cachedRequest;
    });
}

function loadFromStorage() {
    return Rx.Observable.create(ob => {
        var storedData = {  // loading from storage goes here
            value: 15,      // wrapped in observable to delay loading until subscribed
            refreshtoken: 63, // other process may have updated this between requests
            expires: new Date().getTime() - (60 * 1000) // pretend to have already expired
        };

        ob.next(storedData);
        ob.completed();
    })
}

function saveToStorage(data) {
    // save goes here
}

// first request
getOrUpdateValue().subscribe(function(v) { console.log('sub1: ' + v); }); 

// second request, can occur before or after first request finishes
window.setTimeout(
    () => getOrUpdateValue().subscribe(function(v) { console.log('sub2: ' + v); }),
    1500); 

【问题讨论】:

  • 您在第 3 阶段参数的第二部分中提到。除了refreshtoken,还有其他参数吗?他们是从哪里传过来的?它们是如何比较的?
  • 唯一的参数是刷新令牌。我会更新问题以更清楚
  • 谢谢。令牌是强制性的还是您只是使用它来确定是否需要发出新的后端请求?
  • 令牌是强制性的。如果加载的数据为空或刷新令牌为空,则它会快捷方式并返回 observable.empty
  • 没关系!我很欣赏细节。我已经更新以反映到期时间而不是到期标志

标签: javascript angular rxjs


【解决方案1】:

首先,看看一个有效的jsbin example

解决方案与您的初始代码略有不同,我想解释一下原因。需要不断返回本地存储,保存它,保存标志(缓存和令牌)不适合我的反应性、功能性方法。我给出的解决方案的核心是:

var data$ = new Rx.BehaviorSubject(storageMock);
var request$ = new Rx.Subject();
request$.flatMapFirst(loadFromServer).share().startWith(storageMock).subscribe(data$);
data$.subscribe(saveToStorage);

function getOrUpdateValue() {
    return data$.take(1)
      .filter(data => (data && data.refreshtoken))
      .switchMap(data => (data.expires > new Date().getTime() 
                    ? data$.take(1)
                    : (console.log('expired ...'), request$.onNext(true) ,data$.skip(1).take(1))));
}

关键是data$ 保存您的最新数据并且始终是最新的,通过data$.take(1) 可以轻松访问。 take(1) 对于确保您的订阅获得单个值并终止很重要(因为您尝试以程序方式而不是功能方式工作)。如果没有take(1),您的订阅将保持活动状态,并且您将拥有多个处理程序,也就是说,您还将在仅用于当前更新的代码中处理未来的更新。

此外,我持有一个request$ 主题,这是您开始从服务器获取新数据的方式。该函数的工作原理如下:

  1. 过滤器确保如果您的数据为空或没有令牌,则不会有任何东西通过,类似于您拥有的 return Rx.Observable.empty()
  2. 如果数据是最新的,则返回data$.take(1),这是您可以订阅的单元素序列。
  3. 如果没有,则需要刷新。为此,它会触发request$.onNext(true) 并返回data$.skip(1).take(1)skip(1) 是为了避免当前的过时值。

为简洁起见,我使用了(console.log('expired ...'), request$.onNext(true) ,data$.skip(1).take(1)))。这可能看起来有点神秘。它使用 js 逗号分隔的语法,这在 minifiers/uglifiers 中很常见。它执行所有语句并返回最后一条语句的结果。如果你想要一个更易读的代码,你可以这样重写它:

.switchMap(data => {
  if(data.expires > new Date().getTime()){ 
    return data$.take(1);
  } else {
    console.log('expired ...');
    request$.onNext(true);
    return data$.skip(1).take(1);
  }
});

最后一部分是flatMapFirst的用法。这确保了一旦请求正在进行,所有后续请求都将被丢弃。您可以在控制台打印输出中看到它有效。 “从服务器加载”被打印了几次,但实际序列只被调用一次,你得到一个“从服务器加载完成”打印输出。对于您原来的 refreshtoken 标志检查来说,这是一个更加被动的解决方案。

虽然我不需要保存的数据,但它被保存是因为您提到您可能想在以后的会话中阅读它。

关于 rxjs 的一些小技巧:

  1. 您可以简单地使用Rx.Observable.timer(time_out_value).subscribe(...),而不是使用可能导致很多问题的setTimeout

  2. 创建一个 observable 很麻烦(你甚至不得不调用 next(...)complete())。您可以使用Rx.Subject 以更简洁的方式执行此操作。请注意,您有此类的规范,BehaviorSubjectReplaySubject。这些课程值得了解,并且可以提供很多帮助。

最后一点。这是一个相当大的挑战 :-) 我不熟悉您的服务器端代码和设计注意事项,但抑制呼叫的需要让我感到不舒服。除非有与您的后端相关的很好的理由,否则我的自然方法是使用 flatMap 并让最后一个请求“获胜”,即丢弃以前未终止的调用并设置值。

代码基于 rxjs 4(因此它可以在 jsbin 中运行),如果您使用的是 angular2(因此 rxjs 5),则需要对其进行调整。看看the migration guide

================ 回答史蒂夫的其他问题(在下面的 cmets 中)=======

我可以推荐one article。它的标题说明了一切:-)

至于程序与功能的方法,我会在服务中添加另一个变量:

let token$ = data$.pluck('refreshtoken');

然后在需要时使用它。

我的一般方法是首先映射我的数据流和关系,然后像一个优秀的“键盘管道工”(就像我们所有人一样)构建管道。我的顶级服务草稿将是(为简洁起见,跳过 angular2 手续和提供者):

class UserService {
  data$: <as above>;
  token$: data$.pluck('refreshtoken');
  private request$: <as above>;

  refresh(){
    request.onNext(true);
  }
}

您可能需要进行一些检查,以免pluck 失败。

然后,每个需要数据或令牌的组件都可以直接访问它。

现在假设您有一个服务需要对数据或令牌的更改采取行动:

class SomeService {
  constructor(private userSvc: UserService){
    this.userSvc.token$.subscribe(() => this.doMyUpdates());
  }
}

如果您需要合成数据,意思是,使用数据/令牌和一些本地数据:

Rx.Observable.combineLatest(this.userSvc.data$, this.myRelevantData$)
  .subscribe(([data, myData] => this.doMyUpdates(data.someField, myData.someField));

同样,理念是您构建数据流和管道,将它们连接起来,然后您所要做的就是触发东西。

我想出的“迷你模式”是在触发序列后传递给服务并注册到结果。让我们以自动完成为例:

class ACService {
   fetch(text: string): Observable<Array<string>> {
     return http.get(text).map(response => response.json().data;
   }
}

然后你必须在每次文本更改时调用它并将结果分配给你的组件:

<div class="suggestions" *ngFor="let suggestion; of suggestions | async;">
  <div>{{suggestion}}</div>
</div>

在你的组件中:

onTextChange(text) {
  this.suggestions = acSVC.fetch(text);
}

但也可以这样:

class ACService {
   createFetcher(textStream: Observable<string>): Observable<Array<string>> {
     return textStream.flatMap(text => http.get(text))
       .map(response => response.json().data;
   }
}

然后在你的组件中:

textStream: Subject<string> = new Subject<string>();
suggestions: Observable<string>;

constructor(private acSVC: ACService){
  this.suggestions = acSVC.createFetcher(textStream);
}

onTextChange(text) {
  this.textStream.next(text);
}

模板代码保持不变。

这似乎是一件小事,但是一旦应用程序变得更大,并且数据流变得复杂,这会更好。您有一个保存数据的序列,您可以在需要的任何地方在组件周围使用它,甚至可以进一步转换它。例如,假设您需要知道建议的数量,在第一种方法中,一旦您得到结果,您需要进一步查询它才能得到它,因此:

onTextChange(text) {
  this.suggestions = acSVC.fetch(text);
  this.suggestionsCount = suggestions.pluck('length'); // in a sequence
  // or
  this.suggestions.subscribe(suggestions => this.suggestionsCount = suggestions.length); // in a numeric variable.
}

现在在第二种方法中,您只需定义:

constructor(private acSVC: ACService){
  this.suggestions = acSVC.createFetcher(textStream);
  this.suggestionsCount = this.suggestions.pluck('length');
}

希望这会有所帮助:-)

在写作的过程中,我试图反思我是如何使用这样的响应式的。不用说,在进行实验时,大量的 jsbin 和奇怪的失败是其中很大一部分。我认为有助于塑造我的方法的另一件事(尽管我目前没有使用它)是学习 redux 并阅读/尝试一些 ngrx(angular 的 redux 端口)。哲学和方法甚至不允许您进行程序化思考,因此您必须调整基于功能、数据、关系和流程的思维方式。

【讨论】:

  • 哇,谢谢!这一切都非常有用。我可以看到您如何将设计颠倒过来,以实现我只能通过对行为进行微观管理才能做到的事情。我只能将请求发送到服务器一次,因为刷新令牌在收到请求时在服务器端被消耗,并且具有相同令牌的未来请求将返回 400,所以感觉“正确”一次发出客户端请求并拥有休息加入它。
  • 你提到我尝试按程序工作。暴露 getOrUpdateValue() 的服务被其他服务用来获取该值并在其他请求中使用它(实际上,该值是一个访问令牌)。你会在 Rx 中以不同的方式处理这个问题吗?对于进入 Rx 思维模式,您有任何可访问的“推荐阅读”吗?
  • 谢谢 Meir,真的不能要求更多了 :)
  • 很高兴我能帮上忙。记住 rx 很有趣,就像整天解决复杂的难题 ;-) 并且感谢挑战
  • 我已经阅读了您提供的所有材料,现在我的设计更加被动。不过,我有几个问题。 (1) 如果我们在getOrUpdateValue() 中没有在data$take(1) 并且令牌总是过期,这会导致将值推入自身的无限循环吗? (2) skip(1) 是否等到源 (request$) 推送一个新值后再继续? (3) 我认为request$ 上的startsWith() 使subject$ 的默认值变得多余,因为它会立即被替换,我是否正确,或者我是否缺少细微差别?
猜你喜欢
  • 2018-08-08
  • 1970-01-01
  • 1970-01-01
  • 2021-05-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多