Empresas
Empleos
  • Sobre nosotros
  • Soluciones
    • Publicación de vacantes
      Publica tu vacante y recibe candidatos calificados en 48h.
    • Evaluación de candidatos
      500+ pruebas técnicas y psicológicas, más anti-fraude.
    • Headhunting
      Búsqueda ejecutiva a la medida de principio a fin.
    • Nómina + EOR
      Dispersión de nómina y EOR en más de 15 países de LATAM.
  • Precios
  • Empleos

0

261
Vistas
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 Respuestas
Responde la pregunta

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 Denunciar

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 Denunciar

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 Denunciar
Responde la pregunta
Encuentra empleos remotos

¡Descubre la nueva forma de encontrar empleo!

Top de empleos
Top categorías de empleo
Empresas
Publicar vacante Precios Comercial
Legal
Términos y condiciones Política de privacidad
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Recomiéndame algunas ofertas
Necesito ayuda