Subject triggering only once after pipe and mergeMap

Viewed 213

This is the declaration for the subject and observable:

private subscription: Subscription;
private subject = new Subject();
public observable$: Observable<any> = defer(() => this.subject)

I am having an issue with resolving a flow involving the flattening of the following nested subscriptions:

// This works

this.subscription = this.observable$.subscribe(
    (response) => {
        console.log("triggered")
        return this.otherService.function(response).subscribe(
            (response) => {
                // handle if success
            },
            (error) => {
                // handle if fail
            }
        )
    },
);

Now, I've tried using a pipe and a mergeMap/switchMap/concatMap but none of them seem to work. The console.log("triggered") is only called once, despite the subject changing values over time. Exemple:

this.subsctiption = this.observable$.pipe(
    concatMap((response) => {
        console.log("triggered")
        return this.otherService.function(response)
    })
).subscribe(
    (response) => {
        // handle if success
    },
    (error) => {
        // handle if error
    }
)

Finally, this is how the subject is calling next. It's a function handle from an event output from child component

public handleNextValue(nextValue) {
  console.log(nextValue) // this also prints successfully so the value is passed correctly
  this.subject.next(nextValue);
}

What I am doing wrong here? Why the first approach works while the 2nd is only triggering the console.log only once?


Later edit: Thanks to Chris Hamilton there is a stackblitz example which actually works as intended: https://stackblitz.com/edit/angular-ivy-ppxhpg?file=src/app/app.component.ts However this makes me even more confused on what might be wrong inside my code, since the subject's next value is correctly called.

1 Answers

concatMap operator subscribed to the Observable returned from its projection function and keeps re-emiting its notifications until it completes. Any subsequent emission from the source Observable is stacked inside concatMap(). And this is the problem when using concatMap() with Subjects. Subjects don't complete until you call Subject.complete().

So it looks like the problem is that this.otherService.function(response) never completes and that's why console.log("triggered") prints only once.

Related