groupBy.js 3.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263
  1. import { Observable } from '../Observable';
  2. import { innerFrom } from '../observable/innerFrom';
  3. import { Subject } from '../Subject';
  4. import { operate } from '../util/lift';
  5. import { createOperatorSubscriber, OperatorSubscriber } from './OperatorSubscriber';
  6. export function groupBy(keySelector, elementOrOptions, duration, connector) {
  7. return operate(function (source, subscriber) {
  8. var element;
  9. if (!elementOrOptions || typeof elementOrOptions === 'function') {
  10. element = elementOrOptions;
  11. }
  12. else {
  13. (duration = elementOrOptions.duration, element = elementOrOptions.element, connector = elementOrOptions.connector);
  14. }
  15. var groups = new Map();
  16. var notify = function (cb) {
  17. groups.forEach(cb);
  18. cb(subscriber);
  19. };
  20. var handleError = function (err) { return notify(function (consumer) { return consumer.error(err); }); };
  21. var activeGroups = 0;
  22. var teardownAttempted = false;
  23. var groupBySourceSubscriber = new OperatorSubscriber(subscriber, function (value) {
  24. try {
  25. var key_1 = keySelector(value);
  26. var group_1 = groups.get(key_1);
  27. if (!group_1) {
  28. groups.set(key_1, (group_1 = connector ? connector() : new Subject()));
  29. var grouped = createGroupedObservable(key_1, group_1);
  30. subscriber.next(grouped);
  31. if (duration) {
  32. var durationSubscriber_1 = createOperatorSubscriber(group_1, function () {
  33. group_1.complete();
  34. durationSubscriber_1 === null || durationSubscriber_1 === void 0 ? void 0 : durationSubscriber_1.unsubscribe();
  35. }, undefined, undefined, function () { return groups.delete(key_1); });
  36. groupBySourceSubscriber.add(innerFrom(duration(grouped)).subscribe(durationSubscriber_1));
  37. }
  38. }
  39. group_1.next(element ? element(value) : value);
  40. }
  41. catch (err) {
  42. handleError(err);
  43. }
  44. }, function () { return notify(function (consumer) { return consumer.complete(); }); }, handleError, function () { return groups.clear(); }, function () {
  45. teardownAttempted = true;
  46. return activeGroups === 0;
  47. });
  48. source.subscribe(groupBySourceSubscriber);
  49. function createGroupedObservable(key, groupSubject) {
  50. var result = new Observable(function (groupSubscriber) {
  51. activeGroups++;
  52. var innerSub = groupSubject.subscribe(groupSubscriber);
  53. return function () {
  54. innerSub.unsubscribe();
  55. --activeGroups === 0 && teardownAttempted && groupBySourceSubscriber.unsubscribe();
  56. };
  57. });
  58. result.key = key;
  59. return result;
  60. }
  61. });
  62. }
  63. //# sourceMappingURL=groupBy.js.map