RxJs - 如何使 observable 表现得像队列 异步事务旁白:为什么要使用 defer?

如何解决RxJs - 如何使 observable 表现得像队列 异步事务旁白:为什么要使用 defer?

我正在努力实现下一个目标:

private beginTransaction(): Observable() {
  ..
}

private test(): void {
  this.beginTransaction().subscribe((): void => {
    this.commitTransaction();
  });

  this.beginTransaction().subscribe((): void => {
     this.commitTransaction();
  });
}

beginTransaction 可以被并发调用,但应该延迟 observable 直到第一个或只有一个 beginTransaction 完成。

换句话说:任何时候只能进行一笔交易。

我尝试了什么:

private transactionInProgress: boolean = false;
private canBeginTransaction: Subject<void> = new Subject<void>();

private bla3(): void {

  this.beginTransaction().subscribe((): void => {
    console.log('beginTransaction 1');
    this.commitTransaction();
  });

  this.beginTransaction().subscribe((): void => {
    console.log('beginTransaction 2');
    this.commitTransaction();
  });

  this.beginTransaction().subscribe((): void => {
    console.log('beginTransaction 3');
    this.commitTransaction();
  });
}

private commitTransaction(): void {
  this.transactionInProgress = false;
  this.canBeginTransaction.next();
}

private beginTransaction(): Observable<void> {

  if(this.transactionInProgress) {
    return of(undefined)
    .pipe(
      skipUntil(this.canBeginTransaction),tap((): void => {
        console.log('begin transaction');
      })
    );
  }
  this.transactionInProgress = true;
  return of(undefined);
}

解决方法

您所问的问题非常模糊和笼统。毫无疑问,更受限制的场景可能看起来要简单得多。

无论如何,我在这里创建了一个管道,每次只允许订阅一次 transaction(): Observable

可能的样子:

/****
 * Represents what each transaction does. Isn't concerned about
 * order/timing/'transactionInProgress' or anything like that.
 * 
 * Here is a fake transaction that just takes 3-5 seconds to emit
 * the string: `Hello ${name}`
 ****/
function transaction(args): Observable<string> {
  const name = args?.message;
  const duration = 3000 + (Math.random() * 2000);
  return of("Hello").pipe(
    tap(_ => console.log("starting transaction")),switchMap(v => timer(duration).pipe(
      map(_ => `${v} ${name}`)
    )),tap(_ => console.log("Ending transation"))
  );
}

// Track transactions
let currentTransactionId = 0;
// Start transactions
const transactionSubj = new Subject<any>();
// Perform transaction: concatMap ensures we only start a new one if 
// there isn't a current transaction underway
const transaction$ = transactionSubj.pipe(
  concatMap(({id,args}) => transaction(args).pipe(
    map(payload => ({id,payload}))
  )),shareReplay(1)
);

/****
 * Begin a new transaction,we give it an ID since transactions are
 * "hot" and we don't want to return the wrong (earlier) transactions,* just the current one started with this call.
 ****/
function beginTransaction(args): Observable<any> {
  return defer(() => {
    const currentId = currentTransactionId++;
    transactionSubj.next({id: currentId,args});
    return transaction$.pipe(
      first(({id}) => id === currentId),map(({payload}) => payload)
    );
  })
}

// Queue up 3 transactions,each one will wait for the previous 
// one to complete before it will begin.
beginTransaction({message: "Dave"}).subscribe(console.log);
beginTransaction({message: "Tom"}).subscribe(console.log);
beginTransaction({message: "Tim"}).subscribe(console.log);

异步事务

当前设置要求事务是异步的,否则您可能会丢失第一个事务。解决这个问题的方法并不简单,所以我构建了一个订阅的操作符,然后尽快调用一个函数。

这是:

function initialize<T>(fn: () => void): MonoTypeOperatorFunction<T> {
  return s => new Observable(observer => {
    const bindOn = name => observer[name].bind(observer);
    const sub = s.subscribe({
      next: bindOn("next"),error: bindOn("error"),complete: bindOn("complete")
    });
    fn();
    return {
      unsubscribe: () => sub.unsubscribe
    };
  });
}

这里正在使用:

function beginTransaction(args): Observable<any> {
  return defer(() => {
    const currentId = currentTransactionId++;
    return transaction$.pipe(
      initialize(() => transactionSubj.next({id: currentId,args})),first(({id}) => id === currentId),map(({payload}) => payload)
    );
  })
}

旁白:为什么要使用 defer

考虑重写 beginTransaction:

function beginTransaction(args): Observable<any> {
  const currentId = currentTransactionId++;
  return transaction$.pipe(
    initialize(() => transactionSubj.next({id: currentId,map(({payload}) => payload)
  );
}

在这种情况下,ID 在您调用 beginTransaction 时设置。

// The ID is set here,but it won't be used until subscribed
const preppedTransaction = beginTransaction({message: "Dave"});

// 10 seconds later,that ID gets used.
setTimeout(
  () => preppedTransaction.subscribe(console.log),10000
);

如果在没有初始化运算符的情况下调用 transactionSubj.next,那么这个问题会变得更糟,因为 transactionSubj.next 也会在订阅 observable 前 10 秒被调用(你肯定会错过输出)

问题仍然存在:

如果你想订阅同一个 observable 两次怎么办?

const preppedTransaction = beginTransaction({message: "Dave"});
preppedTransaction.subscribe(
  value => console.log("First Subscribe: ",value)
);
preppedTransaction.subscribe(
  value => console.log("Second Subscribe: ",value)
);

我希望输出为:

First Subscribe: Hello Dave
Second Subscribe: Hello Dave

相反,你得到

First Subscribe: Hello Dave
First Subscribe: Hello Dave
Second Subscribe: Hello Dave
Second Subscribe: Hello Dave

因为您在订阅时没有获得新 ID,所以两个订阅共享一个 ID。 defer 通过在订阅前不分配 ID 来解决此问题。这在管理流中的错误时变得非常重要(让您在 observable 出错后重新尝试)。

,

我不确定我是否正确理解了问题,但在我看来,concatMap 是您正在寻找的运算符。

示例如下

const transactionTriggers$ = from([
  't1','t2','t3'
])

function processTransation(trigger: string) {
  console.log(`Start processing transation triggered by ${trigger}`)
  // do whatever needs to be done and then return an Observable
  console.log(`Transation triggered by ${trigger} processing ......`)
  return of(`Transation triggered by ${trigger} processed`)
}

transactionTriggers$.pipe(
  concatMap(trigger => processTransation(trigger)),tap(console.log)
).subscribe()

您基本上从一个事件流开始,其中每个事件都应该触发交​​易的处理。

然后您使用 processTransaction 函数执行处理交易所需的任何操作。 processTransactio 需要返回一个 Observable,该 Observable 在事务处理完毕然后完成时发出处理结果。

然后在管道中,如果需要,您可以使用 tap 对处理结果做进一步的处理。

您可以尝试this stackblitz中的代码。

版权声明:本文内容由互联网用户自发贡献,该文观点与技术仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 dio@foxmail.com 举报,一经查实,本站将立刻删除。

相关推荐


使用本地python环境可以成功执行 import pandas as pd import matplotlib.pyplot as plt # 设置字体 plt.rcParams[&#39;font.sans-serif&#39;] = [&#39;SimHei&#39;] # 能正确显示负号 p
错误1:Request method ‘DELETE‘ not supported 错误还原:controller层有一个接口,访问该接口时报错:Request method ‘DELETE‘ not supported 错误原因:没有接收到前端传入的参数,修改为如下 参考 错误2:cannot r
错误1:启动docker镜像时报错:Error response from daemon: driver failed programming external connectivity on endpoint quirky_allen 解决方法:重启docker -&gt; systemctl r
错误1:private field ‘xxx‘ is never assigned 按Altʾnter快捷键,选择第2项 参考:https://blog.csdn.net/shi_hong_fei_hei/article/details/88814070 错误2:启动时报错,不能找到主启动类 #
报错如下,通过源不能下载,最后警告pip需升级版本 Requirement already satisfied: pip in c:\users\ychen\appdata\local\programs\python\python310\lib\site-packages (22.0.4) Coll
错误1:maven打包报错 错误还原:使用maven打包项目时报错如下 [ERROR] Failed to execute goal org.apache.maven.plugins:maven-resources-plugin:3.2.0:resources (default-resources)
错误1:服务调用时报错 服务消费者模块assess通过openFeign调用服务提供者模块hires 如下为服务提供者模块hires的控制层接口 @RestController @RequestMapping(&quot;/hires&quot;) public class FeignControl
错误1:运行项目后报如下错误 解决方案 报错2:Failed to execute goal org.apache.maven.plugins:maven-compiler-plugin:3.8.1:compile (default-compile) on project sb 解决方案:在pom.
参考 错误原因 过滤器或拦截器在生效时,redisTemplate还没有注入 解决方案:在注入容器时就生效 @Component //项目运行时就注入Spring容器 public class RedisBean { @Resource private RedisTemplate&lt;String
使用vite构建项目报错 C:\Users\ychen\work&gt;npm init @vitejs/app @vitejs/create-app is deprecated, use npm init vite instead C:\Users\ychen\AppData\Local\npm-
参考1 参考2 解决方案 # 点击安装源 协议选择 http:// 路径填写 mirrors.aliyun.com/centos/8.3.2011/BaseOS/x86_64/os URL类型 软件库URL 其他路径 # 版本 7 mirrors.aliyun.com/centos/7/os/x86
报错1 [root@slave1 data_mocker]# kafka-console-consumer.sh --bootstrap-server slave1:9092 --topic topic_db [2023-12-19 18:31:12,770] WARN [Consumer clie
错误1 # 重写数据 hive (edu)&gt; insert overwrite table dwd_trade_cart_add_inc &gt; select data.id, &gt; data.user_id, &gt; data.course_id, &gt; date_format(
错误1 hive (edu)&gt; insert into huanhuan values(1,&#39;haoge&#39;); Query ID = root_20240110071417_fe1517ad-3607-41f4-bdcf-d00b98ac443e Total jobs = 1
报错1:执行到如下就不执行了,没有显示Successfully registered new MBean. [root@slave1 bin]# /usr/local/software/flume-1.9.0/bin/flume-ng agent -n a1 -c /usr/local/softwa
虚拟及没有启动任何服务器查看jps会显示jps,如果没有显示任何东西 [root@slave2 ~]# jps 9647 Jps 解决方案 # 进入/tmp查看 [root@slave1 dfs]# cd /tmp [root@slave1 tmp]# ll 总用量 48 drwxr-xr-x. 2
报错1 hive&gt; show databases; OK Failed with exception java.io.IOException:java.lang.RuntimeException: Error in configuring object Time taken: 0.474 se
报错1 [root@localhost ~]# vim -bash: vim: 未找到命令 安装vim yum -y install vim* # 查看是否安装成功 [root@hadoop01 hadoop]# rpm -qa |grep vim vim-X11-7.4.629-8.el7_9.x
修改hadoop配置 vi /usr/local/software/hadoop-2.9.2/etc/hadoop/yarn-site.xml # 添加如下 &lt;configuration&gt; &lt;property&gt; &lt;name&gt;yarn.nodemanager.res