Problema abstracto: cada vez que una fuente Observable emite un evento, se debe activar una secuencia de llamadas API y servicios Angular. Algunas de esas invocaciones dependen de resultados anteriores.
En mi ejemplo, la fuente Observable startUpload$ desencadena una serie de invocaciones dependientes.
Usando la desestructuración, esto se puede escribir así:
this.startUploadEvent$.pipe( concatMap(event => this.getAuthenticationHeaders(event)), map(({ event, headers }) => this.generateUploadId(event, headers)), tap(({ event, headers, id }) => this.emitUploadStartEvent(id, event)), concatMap(({ event, headers, id }) => this.createPdfDocument(event, headers, id)), concatMap(({ event, headers, id, pdfId }) => this.uploadBilderForPdf(event, pdfId, headers, id)), mergeMap(({ event, headers, id, pdfId, cloudId }) => this.closePdf(cloudId, event, headers, id, pdfId)), tap(({ event, headers, id, pdfId, cloudId }) => this.emitUploadDoneEvent(id, event, cloudId)), ).subscribe()Casi se lee como un enfoque imperativo. Pero tiene ciertos problemas:
{ event, headers, id, pdfId, cloudId }generateUploadId(event, headers) ) deben recibir todos los valores anteriores para que puedan pasarlos a la siguiente canalización, incluso si el método en sí no lo requiere._
private closePdf(cloudId, event, headers, id, pdfId) { return this.httpClient.post(..., { headers } ) .pipe( //..., map(() => ({ event, headers, id, pdfId, cloudId })) ) } Sería bueno si el compilador pudiera encargarse del repetitivo (como con async await ) para escribir el código que se lee así (sin ninguno de los problemas mencionados anteriormente):
private startUpload(event: StartUploadEvent) { const headers = this.getAuthenticationHeaders(event) const id = this.generateUploadId() this.emitUploadStartEvent(id, event) const pdfId = this.createPdfDocument(event, headers, id) this.uploadBilderForPdf(event, pdfId, headers, id) const cloudId = this.closePdf(headers, pdfId) this.emitUploadDoneEvent(id, event, cloudId) return cloudId }¿Cómo pasar resultados entre observables encadenados sin los problemas que he mencionado? ¿Hay algún concepto de rxjs que me haya perdido?
¿Podría usar un objeto para el conjunto de datos? Algo como esto:
Interfaz:
export interface Packet { event: string; headers?: string; id?: number; pdfId?: number; cloudId?: number; }Luego en el código, algo como esto:
Servicio:
this.startUploadEvent$.pipe( concatMap(packet => this.doThingOne(packet)), map(packet => this.doThingTwo(packet)), tap(packet => this.doThingThree(packet)), // ... );De esa manera, cada método puede usar los bits del objeto que necesita y pasar el resto. Aunque esto requiere cambiar cada uno de los métodos para aceptar y trabajar con el objeto.
Tiene razón sobre los problemas que produce dicho código y la solución abstracta es trasladar la responsabilidad de combinar los resultados y pasar los argumentos correctos a cada llamada de los métodos a la canalización.
Algunas mejoras se pueden hacer muy fácilmente. El operador tap no modifica el valor, por lo que puede eliminar las propiedades innecesarias de la desestructuración. map simplemente transforma el resultado, así que en su lugar
map(({ event, headers }) => this.generateUploadId(event, headers)),podemos escribir
map(({ event, headers }) => ({ event, headers, id: this.generateUploadId(event, headers) })) y this.generateUploadId ya no tiene que devolver un objeto.
En cuanto a los operadores de mapeo de alto orden, hay algunas opciones que me vinieron a la mente. En primer lugar, la mayoría de los operadores 'xMap' admiten el selector de resultados como último argumento y su propósito es exactamente lo que necesitamos: combinar el valor de origen con el resultado. Los selectores de resultados quedaron obsoletos , por lo que las canalizaciones anidadas son la forma actual de hacerlo, pero echemos un vistazo a cómo se vería usando el selector de resultados
this.startUploadEvent$ .pipe( concatMap( event => this.getAuthenticationHeaders(event), (event, headers) => ({ event, headers }) // <-- Result Selector ) ); Se ve muy similar a la Opción 0, pero el event se mantiene en el cierre en lugar del observable interno.
this.startUploadEvent$ .pipe( concatMap( event => this.getAuthenticationHeaders(event) .pipe(map(headers => ({ event, headers }))) ) );Es posible crear un operador personalizado y obtener una sintaxis bastante similar a los selectores de resultados
function withResultSelector(operator, transformer) { let sourceValue; return pipe( tap(value => (sourceValue = value)), operator, map(value => transformer(sourceValue, value)) ); }Uso:
this.startUploadEvent$ .pipe( withResultSelector( concatMap(event => this.getAuthenticationHeaders(event)), (event, headers) => ({ event, headers }) ) );Yendo más allá, es posible extraer cosas repetitivas y hacer que todo sea más funcional:
const mergeAs = propName => (a, b) => ({ ...a, [propName]: b }); const opAndMergeAs = (operator, propName) => withResultSelector(operator, mergeAs(propName)); this.startUploadEvent$ .pipe( opAndMergeAs(concatMap(event => this.getAuthenticationHeaders(event)), "headers") );Puede ser un poco engorroso escribir tipos adecuados para eso, pero es un problema diferente
Zona de juegos que usé para escribir la respuesta.
Sus métodos definitivamente no deben estar acoplados al contexto, así como no pensar en asignar el resultado a la forma específica.
RxJS tiene que ver con la programación funcional. Y en la programación funcional hay un patrón como Adaptación de argumentos a parámetros ref .
Nos permite desacoplar la firma de métodos del contexto.
Para lograr esto, puede escribir la versión según el contexto de los operadores map , contentMap , mergMap para que la solución final se vea así:
this.startUploadEvent$.pipe( map(withKey('event')), concatMap_(({event}) => this.getAuthenticationHeaders(event), 'headers'), map_(({ headers }) => this.generateUploadId(headers), 'id'), tap(({ event, id }) => this.emitUploadStartEvent(id, event)), concatMap_(({ id }) => this.createPdfDocument(id), 'pdfId'), concatMap_(({ pdfId }) => this.uploadBuilderForPdf(pdfId), 'cloudId'), mergeMap_(({ cloudId }) => this.closePdf(cloudId)), tap(({id, event, cloudId}) => this.emitUploadDoneEvent(id, event, cloudId)), ).subscribe(console.log); Nota _ después de esos operadores.
El objetivo de esos operadores personalizados es tomar el objeto de parámetros, pasar por la función de proyección y agregar el resultado de la proyección al objeto de parámetros original.
function map_<K extends string, P, V>(project: (params: P) => V): OperatorFunction<P, P>; function map_<K extends string, P, V>(project: (params: P) => V, key: K): OperatorFunction<P, P & Record<K, V>>; function map_<K extends string, P, V>(project: (params: P) => V, key?: K): OperatorFunction<P, P> { return map(gatherParams(project, key)); } function concatMap_<K extends string, P, V>(projection: (params: P) => Observable<V>): OperatorFunction<P, P>; function concatMap_<K extends string, P, V>(projection: (params: P) => Observable<V>, key: K): OperatorFunction<P, P & Record<K, V>>; function concatMap_<K extends string, P, V>(projection: (params: P) => Observable<V>, key?: K): OperatorFunction<P, P> { return concatMap(gatherParamsOperator(projection, key)); } function mergeMap_<K extends string, P, V>(projection: (params: P) => Observable<V>): OperatorFunction<P, P>; function mergeMap_<K extends string, P, V>(projection: (params: P) => Observable<V>, key: K): OperatorFunction<P, P & Record<K, V>>; function mergeMap_<K extends string, P, V>(projection: (params: P) => Observable<V>, key?: K): OperatorFunction<P, P> { return mergeMap(gatherParamsOperator(projection, key)); } // https://github.com/Microsoft/TypeScript/wiki/FAQ#why-am-i-getting-supplied-parameters-do-not-match-any-signature-error function gatherParams<K extends string, P, V>(fn: (params: P) => V): (params: P) => P; function gatherParams<K extends string, P, V>(fn: (params: P) => V, key: K): (params: P) => P & Record<K, V>; function gatherParams<K extends string, P, V>(fn: (params: P) => V, key?: K): (params: P) => P { return (params: P) => { if (typeof key === 'string') { return Object.assign({}, params, { [key]: fn(params) } as Record<K, V>); } return params; }; } function gatherParamsOperator<K extends string, P, V>(fn: (params: P) => Observable<V>): (params: P) => Observable<P>; function gatherParamsOperator<K extends string, P, V>(fn: (params: P) => Observable<V>, key: K): (params: P) => Observable<P & Record<K, V>>; function gatherParamsOperator<K extends string, P, V>(fn: (params: P) => Observable<V>, key?: K): (params: P) => Observable<P> { return (params: P) => { return fn(params).pipe(map(value => gatherParams((_: P) => value, key)(params))); }; } function withKey<K extends string, V>(key: K): (value: V) => Record<K, V> { return (value: V) => ({ [key]: value } as Record<K, V>); } Usé sobrecargas de funciones aquí porque a veces no necesitamos agregar claves adicionales a los parámetros. Los parámetros solo deben pasar a través de él en el caso de this.closePdf(...) .
Como resultado, obtiene una versión desacoplada de la misma que tenía antes con seguridad de tipos:
En la mayoría de los casos, debe seguir el principio YAGNI (no lo necesitará). Y sería mejor no agregar más complejidad al código existente. Para tal escenario, debe ceñirse a alguna implementación simple de compartir parámetros entre operadores de la siguiente manera:
ngOnInit() { const params: Partial<Params> = {}; this.startUploadEvent$.pipe( concatMap(event => (params.event = event) && this.getAuthenticationHeaders(event)), map(headers => (params.headers = headers) && this.generateUploadId(headers)), tap(id => (params.uploadId = id) && this.emitUploadStartEvent(id, event)), concatMap(id => this.createPdfDocument(id)), concatMap(pdfId => (params.pdfId = pdfId) && this.uploadBuilderForPdf(pdfId)), mergeMap(cloudId => (params.cloudId = cloudId) && this.closePdf(cloudId)), tap(() => this.emitUploadDoneEvent(params.pdfId, params.cloudId, params.event)), ).subscribe(() => { console.log(params) }); donde Params es:
interface Params { event: any; headers: any; uploadId: any; pdfId: any; cloudId: any; } Tenga en cuenta los paréntesis que utilicé en las tareas (params.cloudId = cloudId) .
También hay muchos otros métodos, pero requieren cambiar su flujo de uso de operadores rxjs:
Por lo que entendí, le preocupa la legibilidad y no tener que llevar la carga útil de un método a otro.
¿Alguna vez has pensado en convertir un Observable en una Promesa? Lo importante aquí es que los observables deben completarse para que la promesa se cumpla y pueda resolverse (es lo mismo que completar pero solo para promesa).
Debido a su consejo, vea arriba (como con async await) llegué a esta sugerencia.
private async startUpload(event: StartUploadEvent) { const headers = await this.getAuthenticationHeaders(event).toPromise(); const id = await this.generateUploadId().toPromise(); this.emitUploadStartEvent(id, event); const pdfId = await this.createPdfDocument(event, headers, id).toPromise(); await this.uploadBilderForPdf(event, pdfId, headers, id).toPromise(); const cloudId = await this.closePdf(headers, pdfId).toPromise(); this.emitUploadDoneEvent(id, event, cloudId) return cloudId }Información: Aquí puede leer lo que sucede si convierte un observable en una promesa sin haber completado el observable: ¿Por qué la promesa convertida del Sujeto (Observable) no funciona como se esperaba?
Nota: Estoy cumpliendo con sus expectativas de acuerdo
Y tal vez haya otras formas de resolver el problema que no violen las mejores prácticas comunes.
Puede:
asignar el resultado de cada acción a un observable
encadenar llamadas de función posteriores basadas en resultados anteriores
esos resultados se pueden reutilizar en llamadas de acción posteriores a través withLatestFrom
shareReplay se usa para evitar que las últimas suscripciones withLatestFrom provoquen que las funciones anteriores se vuelvan a ejecutar
function startUpload(event$: Observable<string>) { const headers$ = event$.pipe( concatMap(event => getAuthenticationHeaders(event)), shareReplay() ); const id$ = headers$.pipe( map(() => generateUploadId()), shareReplay() ); const emitUploadEvent$ = id$.pipe( withLatestFrom(event$), // use earlier result map(([id, event]) => emitUploadStartEvent(id, event)), shareReplay() ); // etc }Como arriba, las funciones solo toman los parámetros que requieren y no hay traspaso.
Demostración: https://stackblitz.com/edit/so-rxjs-chaining-1?file=index.ts
Este patrón se puede simplificar mediante el uso de un operador personalizado rxjs (tenga en cuenta que esto podría refinarse aún más, incluida la escritura):
function call<T, R, TArgs extends any[], OArgs extends Observable<any>[]>( operator: (func: ((a: TArgs) => R)) => OperatorFunction<TArgs,R>, action: (...args: any[]) => R, ignoreInput: boolean, ...observableArgs: OArgs ): (args: Observable<T>) => Observable<R> { return (input: Observable<T>) => input.pipe( withLatestFrom(...observableArgs), operator((args: any[]) => action(...args.slice(ignoreInput ? 1: 0))), shareReplay(1) ); }Que se puede utilizar como:
function startUpload(event$: Observable<string>) { const headers$ = event$.pipe( call(concatMap, getAuthenticationHeaders, true) ); const id$ = headers$.pipe( call(map, generateUploadId, false) ); const startEmitted$ = id$.pipe( call(map, emitUploadStartEvent, true, event$) ); const pdfId$ = startEmitted$.pipe( call(map, createPdfDocument, false, event$, headers$, id$) ); const uploaded$ = pdfId$.pipe( call(map, uploadBuilderForPdf, false, event$, pdfId$, headers$, id$) ); const cloudId$ = uploaded$.pipe( call(map, closePdf, false, headers$, pdfId$) ); const uploadDone$ = cloudId$.pipe( call(map, emitUploadDoneEvent, true, id$, event$) ); // return cloudId$ instead of uploadDone$ but preserve observable chain return uploadDone$.pipe(concatMap(() => cloudId$)); }Demostración: https://stackblitz.com/edit/so-rxjs-chaining-4?file=index.ts
¡Ciertamente no debería hacer que sus métodos tomen parámetros que no les conciernen!
A tu pregunta principal:
¿Cómo pasar resultados entre observables encadenados sin los problemas que he mencionado?
El siguiente código es equivalente a su código de muestra, sin necesidad de pasar las propiedades innecesarias. Los valores devueltos anteriormente son accesibles mediante llamadas a funciones más abajo en la cadena:
1 startUploadEvent$.pipe( 2 concatMap(event => getAuthenticationHeaders(event).pipe( 3 map(headers => generateUploadId(event, headers).pipe( 4 tap(id => emitUploadStartEvent(id, event)), 5 concatMap(id => createPdfDocument(event, headers, id)), 6 concatMap(pdfId => uploadBilderForPdf(event, pdfId)), 7 tap(cloudId => closePdf(cloudId, event)) 8 )) 9 )) 10 ).subscribe(); Observe cómo los event y los headers son accesibles en sentido descendente. No es necesario pasarlos a funciones que no los requieran.
¿Hay algún concepto de rxjs que me haya perdido?
Quizás.? Realmente no... :-)
El truco consiste en agregar un .pipe para agrupar efectivamente a los operadores para que todos tengan acceso a los parámetros de entrada.
Por lo general, tratamos de mantener el código plano dentro del .pipe :
1 const greeting$ = userId$.pipe( 2 switchMap(id => http.get(`/users/${id}`)), 3 map(response => response.data.userName), 4 map(name => `Hello ${name}!`), 5 tap(greeting => console.log(greeting)) 6 );pero ese código realmente no es diferente a:
1 const greeting$ = userId$.pipe( 2 switchMap(id => http.get(`/users/${id}`).pipe( 3 map(response => response.data.userName), 4 map(name => `Hello ${name}! (aka User #${id})`) 5 )), 6 tap(greeting => console.log(greeting)) 7 ); Pero, en el segundo caso, la línea #4 tiene acceso al name y al id , mientras que en el primer caso solo tiene acceso al name .
Observe que la firma del primero es userId$.pipe(switchMap(), map(), map(), tap())
El segundo es: userId$.pipe(switchMap(), tap()) .
Tiene razón sobre estas preocupaciones y problemas que mencionó, pero el problema que veo aquí es cambiar su mentalidad de un enfoque imperativo a un enfoque Reactivo/Funcional, pero primero revisemos el código imperativo.
private startUpload(event: StartUploadEvent) { const headers = this.getAuthenticationHeaders(event) const id = this.generateUploadId() this.emitUploadStartEvent(id, event) const pdfId = this.createPdfDocument(event, headers, id) this.uploadBilderForPdf(event, pdfId, headers, id) const cloudId = this.closePdf(headers, pdfId) this.emitUploadDoneEvent(id, event, cloudId) return cloudId } Aquí ve que el material está más limpio que tiene un event que puede pasar y obtener solo lo que quiere y pasarlo a las siguientes funciones y queremos mover este código al enfoque Reactivo/Funcional.
el principal problema desde mi punto de vista es que hizo que su función perdiera el contexto que tiene, por ejemplo, getAuthenticationHeaders no debería devolver el event en absoluto, solo debería devolver headers y lo mismo para otras funciones.
cuando se trata de RxJS (también conocido como enfoque reactivo), se trata mucho de estos problemas y eso está bien, ya que mantiene los conceptos funcionales aplicados y hace que su código sea más predecible, ya que los operadores pure solo deben tratar con datos en la misma canalización que mantiene todo puro y no provocar efectos secundarios que conducirán a un código impredecible.
Creo que lo que está buscando se resolverá con nested pipes (esta es la mejor solución en mi opinión)
concatMap(event => this.getAuthenticationHeaders(event).pipe( map(headers => this.generateUploadId(event, headers).pipe()) ))y se usa mucho en algunas bibliotecas de respaldo RxJS como Marble.js
puede usar un enfoque similar al Result Selector :
concatMap(event => this.getAuthenticationHeaders(event).pipe( map(headers => ({ headers, event })) )),o las excelentes otras soluciones que la gente sugirió harán que funcione, pero aún tendrá los mismos problemas que menciona pero con un código más limpio/legible.
También puede convertirlo en un enfoque async/await , pero perderá la reactividad que le proporciona RxJS.
lo que puedo sugerir es tratar de leer más sobre la programación reactiva y cómo mover su mentalidad hacia eso y proporcionaré algunos enlaces aquí que veo que son muy buenos para comenzar y probar algunas bibliotecas que se construyeron sobre RxJS. como CycleJS y recomiendo leer sobre Programación funcional, que también ayudará mucho con estos excelentes libros . Guía en su mayoría adecuada para FP (en javascript) y software de composición .
Recomiendo este gran Talk RxJS Recipes que cambiará tu forma de usar RxJS.
Recursos útiles: