Parece haber una extraña discrepancia con la forma en que se cancela la suscripción de share y shareReplay (con refcount:true).
Considere lo siguiente (puede pegarlo en rxviz.com):
const { interval } = Rx; const { take, shareReplay, share, timeoutWith, startWith , finalize} = RxOperators; const shareReplay$ = interval(2000).pipe( finalize(() => console.log('[finalize] Called shareReplay$')), take(1), shareReplay({refcount:true, bufferSize: 0})); shareReplay$.pipe( timeoutWith(1000, shareReplay$.pipe(startWith('X'))), ) const share$ = interval(2000).pipe( finalize(() => console.log('[finalize] Called on share$')), take(1), share()); share$.pipe( timeoutWith(1000, share$.pipe(startWith('X'))), )La salida de streamReplay$ será -X-0- mientras que la salida de shareStream$ será -X--0. Parece que compartir cancela la suscripción de la fuente antes de que timeoutWith pueda volver a suscribirse, mientras que shareReplay logra mantener la suscripción compartida durante el tiempo suficiente para volver a utilizarla.
Quiero usar esto para agregar un tiempo de espera local en un RPC (mientras mantengo la llamada abierta), y cualquier resuscripción sería desastrosa, así que quiero evitar el riesgo de que uno de estos comportamientos sea un error y se cambie en el futuro.
Podría usar race() y fusionar la llamada rpc con un inicio retrasado (para que ambos se suscriban al mismo tiempo), pero sería más código y operadores.
EDITAR: una posible solución es fusionar dos suscripciones, una a una solicitud compartida y otra a una transmisión retrasada que tarda hasta que se emite la transmisión compartida:
merge(share$, of('still working...').pipe( delay(1000), takeUntil(share$)));De esta manera, la transmisión sombreada se suscribe al mismo tiempo, por lo que no hay un "área gris" cuando un operador se da de baja cuando se suscribe un niño. (Convertiré esto en una respuesta a menos que alguien presente una sugerencia mejor) o pueda explicar las intenciones/diferencias entre compartir y compartirReproducir
El tiempo de ejecución de Javascript no viene con garantías de tiempo de ningún tipo. En un entorno de un solo subproceso, eso tiene sentido (y puede comenzar a hacerlo un poco mejor con los trabajadores web y demás). Si ocurre algo que requiere mucha computación, todo lo que se cubre en esa ventana de tiempo solo espera. Sin embargo, la mayoría de las veces sucede en el orden esperado.
Sin importar,
const hello = (name: string) => () => console.log(`Hello ${name}`); setTimeout(hello("first"), 2000); setTimeout(hello("second"), 2000);En Node, V8 o SpiderMonkey, conserva el orden aquí. Pero ¿qué pasa con esto?
const hello = (name: string) => () => console.log(`Hello ${name}`); setTimeout(hello("first"), 1); setTimeout(hello("second"), 0);Aquí asumirías que el segundo siempre es lo primero porque se supone que sucede una milésima de segundo antes. Si ejecuta esto usando SpiderMonkey, el orden depende de qué tan ocupado esté el ciclo de eventos. Cubren tiempos de espera más cortos, ya que un tiempo de espera de 0 ms toma alrededor de 8 ms en promedio de todos modos.
En JavaScript, es una buena práctica nunca hacer implícitas las dependencias de tiempo.
En el siguiente código, podemos saber razonablemente que esos data no estarán indefinidos cuando llamemos a data.value . Esto se basa implícitamente en el intercalado asíncrono:
let data; setTimeout(() => {data = {value: 5};}, 1000); setTimeout(() => console.log(data.value), 2000);Realmente deberíamos hacer explícita esta dependencia. Ya sea comprobando si los datos no están definidos o reestructurando nuestras llamadas
setTimeout(() => { const data = {value: 5}; setTimeout(() => console.log(data.value), 1000); }, 1000);quiere evitar el riesgo de que uno de estos comportamientos sea un error y se modifique en el futuro.
El riesgo real aquí ni siquiera está en cómo la biblioteca implementa la diferencia. También es un riesgo a nivel de idioma. Depende del intercalado asíncrono de cualquier manera. Corre el riesgo de tener un error que solo aparece una vez cada luna azul y no se puede volver a crear/probar fácilmente, etc.
El operador compartido tiene un ShareConfig ( fuente )
export interface ShareConfig<T> { connector?: () => SubjectLike<T>; resetOnError?: boolean | ((error: any) => Observable<any>); resetOnComplete?: boolean | (() => Observable<any>); resetOnRefCountZero?: boolean | (() => Observable<any>); } Si usa vanilla shareReplay(1) o replay({resetOnRefCountZero: false}) , entonces no está confiando en cómo se ordenan los eventos en JS Event Loop.
Aparte rápido:
Sospecho que no obtendrá muchas respuestas hasta que describa el comportamiento que busca.
use esto para agregar un tiempo de espera local en un RPC (mientras mantiene la llamada abierta)
No estoy seguro de lo que implica mantener una llamada RPC abierta. Parece que desea algo que cancele la solicitud dentro de RxJS pero no cancele el RPC en vuelo. Siento que es una falta de coincidencia de dominio. ¿Por qué no mantener sus observables alineados con la semántica de las llamadas que representan?
Habiendo visto su actualización, sospecho que no necesita un tiempo de espera. Parece que el comportamiento que busca se puede lograr sin cancelar la suscripción al observable que representa su RPC. Yo diría que hacerlo simplifica su lógica y hace que sea más fácil de mantener/ampliar en el futuro.
También evita todas las preocupaciones de entrelazado asíncrono que tenía antes.
Aquí tomaré una conjetura educada sobre el comportamiento que buscas. Parece que quieres seguir:
Si ese es el caso, no necesita share en absoluto. Normalmente, un tiempo de espera significa cancelar una solicitud en curso si tarda demasiado. Tal vez no quiera un tiempo de espera en absoluto, quiere permanecer suscrito a la fuente...
Si ese es el caso, aquí hay una manera de hacer esto:
const source$ = rpc(arg1 arg2); // create a unique token (its memory address). // We'll embed a message inside so it's doing double duty const token = {a: "still working..."}; merge( source$, timer(1000).pipe(mapTo(token)) ).pipe( // take(1), but ignoring the token takeWhile(v => v === token, true), // Unwrap the token, emit the contained message map(v => v === token ? va : v) );Por otro lado, si sabe que el envoltorio observable de su RPC nunca emitirá "sigue funcionando...", entonces no necesita un token único y puede simplificar esto para verificar el valor directamente.
merge( rpc(arg1 arg2), timer(1000).pipe(mapTo("still working...")) ).pipe( takeWhile(v => v === "still working...", true) );