-
Notifications
You must be signed in to change notification settings - Fork 81
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Add
RollingAvg
feature to UpdateBy (#3503)
- Loading branch information
Showing
26 changed files
with
2,374 additions
and
44 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
113 changes: 113 additions & 0 deletions
113
...java/io/deephaven/engine/table/impl/updateby/rollingavg/BigDecimalRollingAvgOperator.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,113 @@ | ||
package io.deephaven.engine.table.impl.updateby.rollingavg; | ||
|
||
import io.deephaven.base.RingBuffer; | ||
import io.deephaven.chunk.Chunk; | ||
import io.deephaven.chunk.ObjectChunk; | ||
import io.deephaven.chunk.attributes.Values; | ||
import io.deephaven.engine.table.MatchPair; | ||
import io.deephaven.engine.table.impl.updateby.UpdateByOperator; | ||
import io.deephaven.engine.table.impl.updateby.internal.BaseObjectUpdateByOperator; | ||
import io.deephaven.engine.table.impl.util.RowRedirection; | ||
import org.jetbrains.annotations.NotNull; | ||
import org.jetbrains.annotations.Nullable; | ||
|
||
import java.math.BigDecimal; | ||
import java.math.MathContext; | ||
|
||
public final class BigDecimalRollingAvgOperator extends BaseObjectUpdateByOperator<BigDecimal> { | ||
private static final int RING_BUFFER_INITIAL_CAPACITY = 128; | ||
@NotNull | ||
private final MathContext mathContext; | ||
|
||
protected class Context extends BaseObjectUpdateByOperator<BigDecimal>.Context { | ||
protected ObjectChunk<BigDecimal, ? extends Values> objectInfluencerValuesChunk; | ||
protected RingBuffer<BigDecimal> objectWindowValues; | ||
|
||
protected Context(final int chunkSize) { | ||
super(chunkSize); | ||
objectWindowValues = new RingBuffer<>(RING_BUFFER_INITIAL_CAPACITY); | ||
} | ||
|
||
@Override | ||
public void close() { | ||
super.close(); | ||
objectWindowValues = null; | ||
} | ||
|
||
|
||
@Override | ||
public void setValuesChunk(@NotNull final Chunk<? extends Values> valuesChunk) { | ||
objectInfluencerValuesChunk = valuesChunk.asObjectChunk(); | ||
} | ||
|
||
@Override | ||
public void push(int pos, int count) { | ||
for (int ii = 0; ii < count; ii++) { | ||
BigDecimal val = objectInfluencerValuesChunk.get(pos + ii); | ||
objectWindowValues.add(val); | ||
|
||
// increase the running sum | ||
if (val != null) { | ||
if (curVal == null) { | ||
curVal = val; | ||
} else { | ||
curVal = curVal.add(val, mathContext); | ||
} | ||
} else { | ||
nullCount++; | ||
} | ||
} | ||
} | ||
|
||
@Override | ||
public void pop(int count) { | ||
for (int ii = 0; ii < count; ii++) { | ||
BigDecimal val = objectWindowValues.remove(); | ||
|
||
// reduce the running sum | ||
if (val != null) { | ||
curVal = curVal.subtract(val, mathContext); | ||
} else { | ||
nullCount--; | ||
|
||
} | ||
} | ||
} | ||
|
||
@Override | ||
public void writeToOutputChunk(int outIdx) { | ||
if (objectWindowValues.size() == nullCount) { | ||
outputValues.set(outIdx, null); | ||
curVal = null; | ||
} else { | ||
final BigDecimal count = new BigDecimal(objectWindowValues.size() - nullCount); | ||
outputValues.set(outIdx, curVal.divide(count, mathContext)); | ||
} | ||
} | ||
|
||
|
||
@Override | ||
public void reset() { | ||
super.reset(); | ||
objectWindowValues.clear(); | ||
} | ||
} | ||
|
||
@NotNull | ||
@Override | ||
public UpdateByOperator.Context makeUpdateContext(final int chunkSize) { | ||
return new Context(chunkSize); | ||
} | ||
|
||
public BigDecimalRollingAvgOperator(@NotNull final MatchPair pair, | ||
@NotNull final String[] affectingColumns, | ||
@Nullable final RowRedirection rowRedirection, | ||
@Nullable final String timestampColumnName, | ||
final long reverseWindowScaleUnits, | ||
final long forwardWindowScaleUnits, | ||
@NotNull final MathContext mathContext) { | ||
super(pair, affectingColumns, rowRedirection, timestampColumnName, reverseWindowScaleUnits, | ||
forwardWindowScaleUnits, true, BigDecimal.class); | ||
this.mathContext = mathContext; | ||
} | ||
} |
Oops, something went wrong.