Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions changelog/unreleased/SOLR-18328-support-std-in-rollup.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
title: "Support 'std' (standard deviation) metric in rollup for streaming expressions"
type: added
authors:
- name: khushjain
links:
- name: SOLR-18328
url: https://issues.apache.org/jira/browse/SOLR-18328
Original file line number Diff line number Diff line change
Expand Up @@ -1448,7 +1448,7 @@ For faster aggregation over low to moderate cardinality fields, the `facet` func
* `StreamExpression` (Mandatory)
* `over`: (Mandatory) A list of fields to group by.
* `metrics`: (Mandatory) The list of metrics to compute.
Currently supported metrics are `sum(col)`, `avg(col)`, `min(col)`, `max(col)`, `count(*)`, `missing(col)`, `countDist(col)`, `per(col, percentile)`.
Currently supported metrics are `sum(col)`, `avg(col)`, `min(col)`, `max(col)`, `count(*)`, `missing(col)`, `countDist(col)`, `per(col, percentile)`, `std(col)`.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

is there any more docs on these indivdual ones or is this mention what the pattern is? Just wondering if folks will know how to use it...

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

no dedicated per-metric docs exists, all streaming metrics are documented this same inline way.

Support for missing(col), countDist(col) and per(col, percentile) are also added by me in the past.


=== rollup Syntax

Expand All @@ -1469,7 +1469,8 @@ rollup(
missing(a_i),
countDist(a_i),
per(a_i, 50),
per(a_f, 75)
per(a_f, 75),
std(a_i)
)
----

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,16 +23,15 @@
import org.apache.solr.client.solrj.io.stream.expr.StreamExpressionParameter;
import org.apache.solr.client.solrj.io.stream.expr.StreamFactory;

/**
* Metric that computes the sample standard deviation of a numeric column over a stream. Consistent
* with the {@code std} streaming evaluator.
*/
public class StdMetric extends Metric {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could you include java docs for this class ?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added docstrings.

// How'd the MeanMetric get to be so mean?
// Maybe it was born with it.
// Maybe it was mayba-mean.
//
// I'll see myself out.

private String columnName;
private double doubleSum;
private long longSum;
private double sum;
private double sumSq;
private long count;

public StdMetric(String columnName) {
Expand Down Expand Up @@ -75,21 +74,45 @@ private void init(String functionName, String columnName, boolean outputLong) {
}

@Override
public void update(Tuple tuple) {}
public void update(Tuple tuple) {
Object o = tuple.get(columnName);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could the object 'o' be made immutable by declaring it as final ?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

it could be added but the current state is consistent with the existing style.

double val;
if (o instanceof Double d) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could the 'if-else' structure be simplified ?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

there definitly are some unusla coding patterns in the streaming code... but we tend to follow them once they exist as there are so many of them!

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@epugh Thanks for the clarification :) Is there a documentation about the code standard adopted by the codebase ?
@KhushJain Feel free to fix the if-else or ignore my suggestion :)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

we have "tidy" and errrorprone tools that we run int he builds. As far as in the streaming code, nothign formal written, just, look at other java classes ;-). Having said that, I do appreciate your reviewing this PR, and there are lots of PR's out there that need a reviewer to go through them!

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed on the quirky instanceof pattern, but keeping it consistent with all its siblings.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@epugh I would be happy to review the PR's :)

val = d;
} else if (o instanceof Float f) {
val = f.doubleValue();
} else if (o instanceof Integer i) {
val = i.doubleValue();
} else if (o instanceof Long l) {
val = l.doubleValue();
} else {
return;
}
++count;
sum += val;
sumSq += val * val;
}

@Override
public Metric newInstance() {
return new MeanMetric(columnName, outputLong);
return new StdMetric(columnName, outputLong);
}

@Override
public String[] getColumns() {
return new String[] {columnName};
}

/** Returns the sample standard deviation of the values seen so far. */
@Override
public Number getValue() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could you include javadocs to understand the purpose of this method ?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added docstrings.

return null;
double std =
count <= 1 ? 0.0d : Math.sqrt(((count * sumSq) - (sum * sum)) / (count * (count - 1.0D)));
if (outputLong) {
return Math.round(std);
} else {
return std;
}
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1333,6 +1333,62 @@ public void testDrillStream() throws Exception {
assertEquals(saf, 18, 0);
}

@Test
public void testRollupStdMetric() throws Exception {
new UpdateRequest()
.add(id, "0", "a_s", "hello0", "a_i", "0", "a_f", "1")
.add(id, "2", "a_s", "hello0", "a_i", "2", "a_f", "2")
.add(id, "3", "a_s", "hello3", "a_i", "3", "a_f", "3")
.add(id, "4", "a_s", "hello4", "a_i", "4", "a_f", "4")
.add(id, "1", "a_s", "hello0", "a_i", "1", "a_f", "5")
.add(id, "5", "a_s", "hello3", "a_i", "10", "a_f", "6")
.add(id, "6", "a_s", "hello4", "a_i", "11", "a_f", "7")
.add(id, "7", "a_s", "hello3", "a_i", "12", "a_f", "8")
.add(id, "8", "a_s", "hello3", "a_i", "13", "a_f", "9")
.add(id, "9", "a_s", "hello0", "a_i", "14", "a_f", "10")
.commit(cluster.getSolrClient(), COLLECTIONORALIAS);

ModifiableSolrParams paramsLoc = new ModifiableSolrParams();
String expr =
"rollup("
+ " search(collection1, q=*:*, fl=\"a_s,a_i,a_f\", sort=\"a_s asc\", qt=\"/export\"),"
+ " over=\"a_s\", std(a_i), std(a_f), count(*)"
+ ")";
paramsLoc.set("expr", expr);
paramsLoc.set("qt", "/stream");

String url =
cluster.getJettySolrRunners().get(0).getBaseUrl().toString() + "/" + COLLECTIONORALIAS;
TupleStream solrStream = new SolrStream(url, paramsLoc);

StreamContext context = new StreamContext();
solrStream.setStreamContext(context);
List<Tuple> tuples = getTuples(solrStream);

assertEquals(3, tuples.size());

// hello0: a_i = [0, 1, 2, 14], a_f = [1, 2, 5, 10]
Tuple tuple = tuples.get(0);
assertEquals("hello0", tuple.getString("a_s"));
assertEquals(6.5511, tuple.getDouble("std(a_i)"), 0.001);
assertEquals(4.0415, tuple.getDouble("std(a_f)"), 0.001);
assertEquals(4, tuple.getDouble("count(*)"), 0.0);

// hello3: a_i = [3, 10, 12, 13], a_f = [3, 6, 8, 9]
tuple = tuples.get(1);
assertEquals("hello3", tuple.getString("a_s"));
assertEquals(4.5092, tuple.getDouble("std(a_i)"), 0.001);
assertEquals(2.6458, tuple.getDouble("std(a_f)"), 0.001);
assertEquals(4, tuple.getDouble("count(*)"), 0.0);

// hello4: a_i = [4, 11], a_f = [4, 7]
tuple = tuples.get(2);
assertEquals("hello4", tuple.getString("a_s"));
assertEquals(4.9497, tuple.getDouble("std(a_i)"), 0.001);
assertEquals(2.1213, tuple.getDouble("std(a_f)"), 0.001);
assertEquals(2, tuple.getDouble("count(*)"), 0.0);
}

@Test
public void testFacetStream() throws Exception {

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@
import org.apache.solr.client.solrj.io.stream.metrics.MinMetric;
import org.apache.solr.client.solrj.io.stream.metrics.MissingMetric;
import org.apache.solr.client.solrj.io.stream.metrics.PercentileMetric;
import org.apache.solr.client.solrj.io.stream.metrics.StdMetric;
import org.apache.solr.client.solrj.io.stream.metrics.SumMetric;
import org.apache.solr.client.solrj.request.CollectionAdminRequest;
import org.apache.solr.client.solrj.request.UpdateRequest;
Expand Down Expand Up @@ -1676,6 +1677,8 @@ public void testRollupStream() throws Exception {
new MaxMetric("a_f"),
new MeanMetric("a_i"),
new MeanMetric("a_f"),
new StdMetric("a_i"),
new StdMetric("a_f"),
new CountMetric(),
new MissingMetric("b_f"),
new CountDistinctMetric("a_i"),
Expand All @@ -1700,6 +1703,8 @@ public void testRollupStream() throws Exception {
Double maxf = tuple.getDouble("max(a_f)");
Double avgi = tuple.getDouble("avg(a_i)");
Double avgf = tuple.getDouble("avg(a_f)");
Double stdi = tuple.getDouble("std(a_i)");
Double stdf = tuple.getDouble("std(a_f)");
Double count = tuple.getDouble("count(*)");
Double missingBf = tuple.getDouble("missing(b_f)");
Double countDistI = tuple.getDouble("countDist(a_i)");
Expand All @@ -1714,6 +1719,8 @@ public void testRollupStream() throws Exception {
assertEquals(10, maxf, 0.001);
assertEquals(4.25, avgi, 0.001);
assertEquals(4.5, avgf, 0.001);
assertEquals(6.5511, stdi, 0.001);
assertEquals(4.0415, stdf, 0.001);
assertEquals(4, count, 0.001);
assertEquals(2, missingBf, 0.001);
assertEquals(4, countDistI, 0.001);
Expand All @@ -1729,6 +1736,8 @@ public void testRollupStream() throws Exception {
maxf = tuple.getDouble("max(a_f)");
avgi = tuple.getDouble("avg(a_i)");
avgf = tuple.getDouble("avg(a_f)");
stdi = tuple.getDouble("std(a_i)");
stdf = tuple.getDouble("std(a_f)");
count = tuple.getDouble("count(*)");
missingBf = tuple.getDouble("missing(b_f)");
countDistI = tuple.getDouble("countDist(a_i)");
Expand All @@ -1743,6 +1752,8 @@ public void testRollupStream() throws Exception {
assertEquals(9, maxf, 0.001);
assertEquals(9.5, avgi, 0.001);
assertEquals(6.5, avgf, 0.001);
assertEquals(4.5092, stdi, 0.001);
assertEquals(2.6458, stdf, 0.001);
assertEquals(4, count, 0.001);
assertEquals(3, missingBf, 0.001);
assertEquals(4, countDistI, 0.001);
Expand All @@ -1758,6 +1769,8 @@ public void testRollupStream() throws Exception {
maxf = tuple.getDouble("max(a_f)");
avgi = tuple.getDouble("avg(a_i)");
avgf = tuple.getDouble("avg(a_f)");
stdi = tuple.getDouble("std(a_i)");
stdf = tuple.getDouble("std(a_f)");
count = tuple.getDouble("count(*)");
missingBf = tuple.getDouble("missing(b_f)");
countDistI = tuple.getDouble("countDist(a_i)");
Expand All @@ -1772,6 +1785,8 @@ public void testRollupStream() throws Exception {
assertEquals(7, maxf, 0.01);
assertEquals(7.5, avgi, 0.01);
assertEquals(5.5, avgf, 0.01);
assertEquals(4.9497, stdi, 0.01);
assertEquals(2.1213, stdf, 0.01);
assertEquals(2, count, 0.01);
assertEquals(0, missingBf, 0.01);
assertEquals(2, countDistI, 0.01);
Expand Down