【问题标题】:Rxjs subscription queueRxjs 订阅队列
【发布时间】:2018-09-09 16:17:49
【问题描述】:

我的 Angular 应用中有一个 Firebase 订阅,它会触发多次。 我如何实现将任务作为队列处理,以便我可以将每个任务同步运行一次?

this.tasks.subscribe(async tasks => {
   for (const x of tasks) 
      await dolongtask(x); // has to be sync
      await removetask(x);
   });

问题是当 longtask 仍在处理时会触发 subribe 事件。

【问题讨论】:

  • 使 longtask 返回一个在任务完成时完成的可观察对象并使用 concatMap。请注意,这种背压可能会导致内存泄漏。
  • 为什么我必须从 longtask 返回一个 observable?
  • 因为 concatMap 需要 observables。你当然可以在没有 observables 的情况下做到这一点,但是使用 rxjs 有什么意义呢?
  • 你能粘贴一些代码吗?
  • this.tasks.pipe(concatMap(tasks => processTasks(tasks))).subscribe()。不能说更多,因为你没有说你的功能是做什么的。

标签: javascript firebase rxjs reactive-programming


【解决方案1】:

恕我直言,我会尝试利用 rxjs 的强大功能,因为无论如何我们已经在这里使用它,并且避免按照另一个答案的建议实施自定义排队概念(尽管您当然可以这样做)。

如果我们稍微简化一下给定的情况,我们只是有一些可观察的,并希望为每个发射执行一个长时间运行的过程 - 按顺序。 rxjs 允许通过 concatMap 操作符来实现这一点,基本上是开箱即用的:

$data.pipe(concatMap(item => processItem(item))).subscribe();

这仅假设processItem 返回一个可观察对象。由于您使用了await,我假设您的函数当前返回 Promises。这些可以使用from 轻松转换为可观察对象。

从 OP 剩下的唯一细节是 observable 实际上发出了一个项目数组,我们希望对每个发射的每个项目执行操作。为此,我们只需使用 mergeMap 将 observable 展平即可。


让我们把它们放在一起。请注意,如果您不准备一些存根数据和日志记录,那么实际的实现只需 行代码(使用 mergeMap + concatMap)。

const { from, interval } = rxjs;
const { mergeMap, concatMap, take, bufferCount, tap } = rxjs.operators;

// Stub for the long-running operation
function processTask(task) {
  console.log("Processing task: ", task);
  return new Promise(resolve => {
    setTimeout(() => {
      console.log("Finished task: ", task);
      resolve(task);
    }, 500 * Math.random() + 300);
  });
}

// Turn processTask into a function returning an observable
const processTask$ = item => from(processTask(item));

// Some stubbed data stream
const tasks$ = interval(250).pipe(
  take(9),
  bufferCount(3),
);

tasks$.pipe(
  tap(task => console.log("Received task: ", task)),
  // Flatten the tasks array since we want to work in sequence anyway
  mergeMap(tasks => tasks),
  // Process each task, but do so consecutively
  concatMap(task => processTask$(task)),
).subscribe(null, null, () => console.log("Done"));
<script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/6.3.2/rxjs.umd.js"></script>

【讨论】:

  • 我不明白除了tasks$.subscribe(tasks => tasks.forEach(task => processTask(task)) 之外还有什么作用。
  • 如果您想在mergeMap()concatMap() 之间添加额外的步骤,这看起来像是一个有趣的模式,例如filter(task => isNewTask(task))
  • @Hiram K. Hackenbacker 首先,它回答了问题,这意味着它按顺序处理事物。像你一样调用函数并不能保证这一点。
  • concatMap 运算符中(有效)内部化的队列的唯一问题是您上面提到的背压风险。也许它可以用map(task => delayTaskIfTooMany(task)) 处理,但我还看不到如何将延迟添加到管道中。
  • 顺便说一句,我认为它可以在没有 const processTask$ = item => from(processTask(item)) 的情况下工作,因为 processTask 返回一个承诺(请参阅 Example 2)。
【解决方案2】:

我根据您提供的代码做了几个假设,

  • 其他应用程序将任务添加到 firebase db(异步),并且此代码正在实现任务处理器。

  • 您的 firebase 查询会返回所有未处理的任务(在一个集合中),并且每次添加新任务时都会发出完整列表。

  • 只有在 removeTask() 运行后,查询才会删除任务

如果是这样,您需要在处理器之前使用重复数据删除机制。

为了便于说明,我使用主题(将其重命名为 tasksQuery$)模拟了 firebase 查询,并在脚本底部模拟了一系列 firebase 事件。 我希望它不会太混乱!

console.clear()
const { mergeMap, filter } = rxjs.operators;

// Simulate tasks query  
const tasksQuery$ = new rxjs.Subject();

// Simulate dolongtask and removetask (assume both return promises that can be awaited)
const dolongtask = (task) => {
  console.log( `Processing: ${task.id}`);
  return new Promise(resolve => {
    setTimeout(() => {
      console.log( `Processed: ${task.id}`);
      resolve('done')
    }, 1000);
  });
}
const removeTask = (task) => {
  console.log( `Removing: ${task.id}`);
  return new Promise(resolve => {
    setTimeout(() => {
      console.log( `Removed: ${task.id}`);
      resolve('done')
    }, 200);
  });
}

// Set up queue (this block could be a class in Typescript)
let tasks = [];
const queue$ = new rxjs.Subject();
const addToQueue = (task) => {
  tasks = [...tasks, task];
  queue$.next(task);
}
const removeFromQueue = () => tasks = tasks.slice(1);
const queueContains = (task) => tasks.map(t => t.id).includes(task.id)

// Dedupe and enqueue
tasksQuery$.pipe(
  mergeMap(tasks => tasks), // flatten the incoming task array 
  filter(task => task && !queueContains(task)) // check not in queue
).subscribe(task => addToQueue(task) );

//Process the queue
queue$.subscribe(async task => {
  await dolongtask(task);
  await removeTask(task); // Assume this sends 'delete' to firebase
  removeFromQueue();
});

// Run simulation
tasksQuery$.next([{id:1},{id:2}]);
// Add after delay to show repeated items in firebase
setTimeout(() => {
  tasksQuery$.next([{id:1},{id:2},{id:3}]); 
}, 500);
<script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/6.3.2/rxjs.umd.js"></script>

【讨论】:

    【解决方案3】:

    撇开你的标题“Rxjs 订阅队列”不谈,你实际上可以修复你的 async/await 代码。

    问题是 async/await 不能很好地与 for 循环配合使用,请参阅这个问题 Using async/await with a forEach loop

    例如,您可以按照@Bergi 的回答替换 for 循环,

    Promise.all()

    console.clear();
    const { interval } = rxjs;
    const { take, bufferCount } = rxjs.operators;
    
    function processTask(task) {
      console.log(`Processing task ${task}`);
      return new Promise(resolve => {
        setTimeout(() => {
          resolve(task);
        }, 500 * Math.random() + 300);
      });
    }
    function removeTask(task) {
      console.log(`Removing task ${task}`);
      return new Promise(resolve => {
        setTimeout(() => {
          resolve(task);
        }, 50);
      });
    }
    
    const tasks$ = interval(250).pipe(
      take(10),
      bufferCount(3),
    );
    
    tasks$.subscribe(async tasks => {
      await Promise.all(
        tasks.map(async task => {
          await processTask(task); // has to be sync
          await removeTask(task);
          console.log(`Finished task ${task}`);
        })
      );
    });
    <script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/6.3.2/rxjs.umd.js"></script>

    更好的是,您可以调整查询以避免使用 for 循环,

    mergeMap()

    console.clear();
    const { interval } = rxjs;
    const { mergeMap, take, bufferCount } = rxjs.operators;
    
    function processTask(task) {
      console.log(`Processing task ${task}`);
      return new Promise(resolve => {
        setTimeout(() => {
          resolve(task);
        }, 500 * Math.random() + 300);
      });
    }
    function removeTask(task) {
      console.log(`Removing task ${task}`);
      return new Promise(resolve => {
        setTimeout(() => {
          resolve(task);
        }, 50);
      });
    }
    
    const tasks$ = interval(250).pipe(
      take(10),
      bufferCount(3),
    );
    
    tasks$
    .pipe(mergeMap(tasks => tasks))
    .subscribe(
      async task => {
        await processTask(task); // has to be sync
        await removeTask(task);
        console.log(`Finished task ${task}`);
      }
    );
    <script src="https://cdnjs.cloudflare.com/ajax/libs/rxjs/6.3.2/rxjs.umd.js"></script>

    【讨论】:

    • forEach 循环不能很好地与 async/await 配合使用,for - 很好。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-12-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多