bufferCount.js 4.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129
  1. import { Subscriber } from '../Subscriber';
  2. /**
  3. * Buffers the source Observable values until the size hits the maximum
  4. * `bufferSize` given.
  5. *
  6. * <span class="informal">Collects values from the past as an array, and emits
  7. * that array only when its size reaches `bufferSize`.</span>
  8. *
  9. * <img src="./img/bufferCount.png" width="100%">
  10. *
  11. * Buffers a number of values from the source Observable by `bufferSize` then
  12. * emits the buffer and clears it, and starts a new buffer each
  13. * `startBufferEvery` values. If `startBufferEvery` is not provided or is
  14. * `null`, then new buffers are started immediately at the start of the source
  15. * and when each buffer closes and is emitted.
  16. *
  17. * @example <caption>Emit the last two click events as an array</caption>
  18. * var clicks = Rx.Observable.fromEvent(document, 'click');
  19. * var buffered = clicks.bufferCount(2);
  20. * buffered.subscribe(x => console.log(x));
  21. *
  22. * @example <caption>On every click, emit the last two click events as an array</caption>
  23. * var clicks = Rx.Observable.fromEvent(document, 'click');
  24. * var buffered = clicks.bufferCount(2, 1);
  25. * buffered.subscribe(x => console.log(x));
  26. *
  27. * @see {@link buffer}
  28. * @see {@link bufferTime}
  29. * @see {@link bufferToggle}
  30. * @see {@link bufferWhen}
  31. * @see {@link pairwise}
  32. * @see {@link windowCount}
  33. *
  34. * @param {number} bufferSize The maximum size of the buffer emitted.
  35. * @param {number} [startBufferEvery] Interval at which to start a new buffer.
  36. * For example if `startBufferEvery` is `2`, then a new buffer will be started
  37. * on every other value from the source. A new buffer is started at the
  38. * beginning of the source by default.
  39. * @return {Observable<T[]>} An Observable of arrays of buffered values.
  40. * @method bufferCount
  41. * @owner Observable
  42. */
  43. export function bufferCount(bufferSize, startBufferEvery = null) {
  44. return function bufferCountOperatorFunction(source) {
  45. return source.lift(new BufferCountOperator(bufferSize, startBufferEvery));
  46. };
  47. }
  48. class BufferCountOperator {
  49. constructor(bufferSize, startBufferEvery) {
  50. this.bufferSize = bufferSize;
  51. this.startBufferEvery = startBufferEvery;
  52. if (!startBufferEvery || bufferSize === startBufferEvery) {
  53. this.subscriberClass = BufferCountSubscriber;
  54. }
  55. else {
  56. this.subscriberClass = BufferSkipCountSubscriber;
  57. }
  58. }
  59. call(subscriber, source) {
  60. return source.subscribe(new this.subscriberClass(subscriber, this.bufferSize, this.startBufferEvery));
  61. }
  62. }
  63. /**
  64. * We need this JSDoc comment for affecting ESDoc.
  65. * @ignore
  66. * @extends {Ignored}
  67. */
  68. class BufferCountSubscriber extends Subscriber {
  69. constructor(destination, bufferSize) {
  70. super(destination);
  71. this.bufferSize = bufferSize;
  72. this.buffer = [];
  73. }
  74. _next(value) {
  75. const buffer = this.buffer;
  76. buffer.push(value);
  77. if (buffer.length == this.bufferSize) {
  78. this.destination.next(buffer);
  79. this.buffer = [];
  80. }
  81. }
  82. _complete() {
  83. const buffer = this.buffer;
  84. if (buffer.length > 0) {
  85. this.destination.next(buffer);
  86. }
  87. super._complete();
  88. }
  89. }
  90. /**
  91. * We need this JSDoc comment for affecting ESDoc.
  92. * @ignore
  93. * @extends {Ignored}
  94. */
  95. class BufferSkipCountSubscriber extends Subscriber {
  96. constructor(destination, bufferSize, startBufferEvery) {
  97. super(destination);
  98. this.bufferSize = bufferSize;
  99. this.startBufferEvery = startBufferEvery;
  100. this.buffers = [];
  101. this.count = 0;
  102. }
  103. _next(value) {
  104. const { bufferSize, startBufferEvery, buffers, count } = this;
  105. this.count++;
  106. if (count % startBufferEvery === 0) {
  107. buffers.push([]);
  108. }
  109. for (let i = buffers.length; i--;) {
  110. const buffer = buffers[i];
  111. buffer.push(value);
  112. if (buffer.length === bufferSize) {
  113. buffers.splice(i, 1);
  114. this.destination.next(buffer);
  115. }
  116. }
  117. }
  118. _complete() {
  119. const { buffers, destination } = this;
  120. while (buffers.length > 0) {
  121. let buffer = buffers.shift();
  122. if (buffer.length > 0) {
  123. destination.next(buffer);
  124. }
  125. }
  126. super._complete();
  127. }
  128. }
  129. //# sourceMappingURL=bufferCount.js.map