Business
Jobs
  • About Us
  • Solutions
    • Job Postings
      Post your job and receive qualified candidates in 48h.
    • Candidate Assessments
      500+ technical and psychological tests, plus anti-fraud.
    • Headhunting
      Tailor-made executive search from start to finish.
    • Payroll + EOR
      Payroll dispersal and EOR across 15+ LATAM countries.
  • Pricing
  • Jobs

0

262
Views
Retraso secuencial RxJS entre emisiones

Tengo un Observable observableA . ¿Cómo puedo crear otro Observable observableB que emita valor de observableA , pero el tiempo entre cada emisión es de al menos 1000 milisegundos?
Por ejemplo:

 observableA: X...500ms..Y..1500ms..Z...600ms..T observableB: X..1000ms..Y..1500ms..Z..1000ms..T
about 4 years ago · Juan Pablo Isaza
3 answers
Answer question

0

Siempre puede usar un temporizador silencioso como límite inferior para cada nueva duración de emisión. Entonces concatMap manejará la contrapresión por usted.

 const bound = timer(1000).pipe(ignoreElements()); const observableB = observableA.pipe( concatMap(v => merge(of(v), bound)) );
about 4 years ago · Juan Pablo Isaza Report

0

Puede usar intervalo y zip para crear un observable que se active 1000 ms después de cada valor en observableA, y luego comprimirlo con observableA:

 const observableA = create(subscriber => { window.setTimeout(() => { subscriber.next(1); window.setTimeout(() => { subscriber.next(2); window.setTimeout(() => { subscriber.next(3); subscriber.complete(); }, 500); }, 1500); }, 600); }); const delay = concat(from([undefined]), observableA).pipe( flatMap(() => interval(1000).take(1)) ); const observableB = zip(observableA, delay).map(([first]) => first); observableB.subscribe((val) => { console.log(JSON.stringify(val) + " " + Date.now()); });

Si no desea retrasar el primer valor, puede hacerlo

 const delay = concat(from([undefined]), observableA.pipe( flatMap(() => interval(1000).take(1)) ));

Para agregar un primer valor inmediato.

about 4 years ago · Juan Pablo Isaza Report

0

Si desea ponerlo en una sola cadena, puede usar lo siguiente usando el operador connect() :

 const obsA$ = range(5).pipe( concatMap((v) => of(v).pipe(delay(Math.random() * 2000))) ); const MIN_DELAY = 1000; let observerReceivedTime = new Date().getTime(); obsA$ .pipe( startWith(null), connect((shared$) => { let lastEmissionTime = null; return shared$.pipe( concatMap((value) => { let delayedValue$; if (!lastEmissionTime) { // Don't delay the first emission. delayedValue$ = of(value); } else { const delayBetweenValues = new Date().getTime() - lastEmissionTime; // Wait for at least `MIN_DELAY`. delayedValue$ = delayBetweenValues > MIN_DELAY ? of(value) : of(value).pipe(delay(MIN_DELAY - delayBetweenValues)); } return delayedValue$.pipe( // Remember time when the last value was reemited. tap(() => (lastEmissionTime = new Date().getTime())), ); }) ); }) ) .subscribe((v) => { const now = new Date().getTime(); console.log( `observer receives emission with at least ${MIN_DELAY}ms delay:`, now - observerReceivedTime, 'value:', v ); observerReceivedTime = now; });

Demostración en vivo: https://stackblitz.com/edit/rxjs-wmfetc?devtoolsheight=60&file=index.ts

Estoy usando el operador RxJS 7 connect() para crear una variable con ámbito lastEmissionTime , por lo que no necesito generar efectos secundarios (aunque, para ser honesto, la cadena sería un poco más corta y más fácil de entender).

Básicamente, solo mide el tiempo entre la emisión anterior y la actual y ajustará delay() para que siempre sea al menos MIN_DELAY .

En la devolución de llamada del observador, solo mido las marcas de tiempo para demostrar que el retraso siempre es de al menos 1000 ms como querías.

Editar: Esta es una alternativa usando mergeScan que es mejor en mi opinión:

 obsA$ .pipe( startWith(null), mergeScan( ([_, lastEmissionTime], currValue) => { let delayedValue$; if (!lastEmissionTime) { // Don't delay the first emission. delayedValue$ = of(currValue); } else { const delayBetweenValues = new Date().getTime() - lastEmissionTime; // Wait for at least `MIN_DELAY`. delayedValue$ = delayBetweenValues > MIN_DELAY ? of(currValue) : of(currValue).pipe(delay(MIN_DELAY - delayBetweenValues)); } return delayedValue$.pipe( // Remember time when the last value was reemited. map((value) => [value, new Date().getTime()]) ); }, [null, null], 1 // We need to set concurrency to `1`. ), map(([value]) => value) )

Demostración en vivo: https://stackblitz.com/edit/rxjs-taeyzx?devtoolsheight=60&file=index.ts

about 4 years ago · Juan Pablo Isaza Report
Answer question
Find remote jobs

Discover the new way to find a job!

Top jobs
Top job categories
Business
Post vacancy Pricing Sales
Legal
Terms and conditions Privacy policy
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Show me some job opportunities
There's an error!