Skip to content
12 changes: 10 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,13 @@ cat known-users.txt | bf build --expected-items 1000000 --false-positive-rate 0.
cat candidates.txt | bf test users.bf
```

Count or rank a tab-separated key without an extra extraction step:

```sh
hll count -d $'\t' -f 2 events.tsv
heavy --top 20 -d $'\t' -f 2 events.tsv
```

Profile JSON field cardinality:

```sh
Expand All @@ -106,10 +113,11 @@ Find frequent API paths:
awk '{print $7}' access.log | heavy --top 20
```

Take a stable sample:
Take a stable sample and its complementary cohort:

```sh
sample --rate 0.01 --stable events.jsonl
sample --rate 0.01 --stable events.jsonl > selected.jsonl
sample --rate 0.01 --stable --invert events.jsonl > remaining.jsonl
```

## Behavior At A Glance
Expand Down
40 changes: 6 additions & 34 deletions cmd/bf/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,6 @@ package main

import (
"bufio"
"bytes"
"encoding/json"
"flag"
"fmt"
Expand Down Expand Up @@ -88,8 +87,6 @@ func build(args []string, in io.Reader, out, errOut io.Writer) error {
noSizeLimit := addSizeLimitFlag(fs)
var input prob.InputOptions
prob.AddInputFlags(fs, &input)
var fields prob.FieldOptions
prob.AddFieldFlags(fs, &fields)
fs.Usage = func() {
fmt.Fprint(fs.Output(), `Usage: bf build --expected-items <n> --false-positive-rate <p> [--delimiter value --field n] [--no-size-limit] [file...] > filter.bf

Expand All @@ -104,15 +101,15 @@ Options:
if err := fs.Parse(args); err != nil {
return err
}
if err := fields.Validate(); err != nil {
if err := input.Fields.Validate(); err != nil {
return fmt.Errorf("usage: %w", err)
}

f, err := bf.NewWithLimit(*expected, *rate, filterSizeLimit(*noSizeLimit))
if err != nil {
return err
}
if err := eachSelectedInput(fs.Args(), in, input, fields, func(item, _ []byte) error {
if err := prob.EachSelectedInputFrom(fs.Args(), in, input, func(item, _ []byte) error {
f.Add(item)
return nil
}); err != nil {
Expand All @@ -130,8 +127,6 @@ func test(args []string, in io.Reader, out io.Writer) error {
noSizeLimit := addSizeLimitFlag(fs)
var input prob.InputOptions
prob.AddInputFlags(fs, &input)
var fields prob.FieldOptions
prob.AddFieldFlags(fs, &fields)
fs.Usage = func() {
fmt.Fprint(fs.Output(), `Usage: bf test [--invert] [--delimiter value --field n] [--no-size-limit] <filter.bf> [file...]

Expand All @@ -147,7 +142,7 @@ Options:
if err := fs.Parse(args); err != nil {
return err
}
if err := fields.Validate(); err != nil {
if err := input.Fields.Validate(); err != nil {
return fmt.Errorf("usage: %w", err)
}
if fs.NArg() < 1 {
Expand All @@ -167,7 +162,7 @@ Options:
return err
}
buffered := bufio.NewWriter(out)
err = eachSelectedInput(paths, in, input, fields, func(item, output []byte) error {
err = prob.EachSelectedInputFrom(paths, in, input, func(item, output []byte) error {
present := f.Test(item)
if present != *invert {
return writeItem(buffered, output, input.NUL)
Expand All @@ -185,8 +180,6 @@ func dedupe(args []string, in io.Reader, out, errOut io.Writer) error {
noSizeLimit := addSizeLimitFlag(fs)
var input prob.InputOptions
prob.AddInputFlags(fs, &input)
var fields prob.FieldOptions
prob.AddFieldFlags(fs, &fields)
fs.Usage = func() {
fmt.Fprint(fs.Output(), `Usage: bf dedupe --expected-items <n> --false-positive-rate <p> [--delimiter value --field n] [--no-size-limit] [file...]

Expand All @@ -201,15 +194,15 @@ Options:
if err := fs.Parse(args); err != nil {
return err
}
if err := fields.Validate(); err != nil {
if err := input.Fields.Validate(); err != nil {
return fmt.Errorf("usage: %w", err)
}
f, err := bf.NewWithLimit(*expected, *rate, filterSizeLimit(*noSizeLimit))
if err != nil {
return err
}
buffered := bufio.NewWriter(out)
err = eachSelectedInput(fs.Args(), in, input, fields, func(item, output []byte) error {
err = prob.EachSelectedInputFrom(fs.Args(), in, input, func(item, output []byte) error {
if f.Test(item) {
return nil
}
Expand Down Expand Up @@ -358,27 +351,6 @@ func countStdin(paths []string) int {
return count
}

func eachSelectedInput(paths []string, in io.Reader, input prob.InputOptions, fields prob.FieldOptions, fn func(item, output []byte) error) error {
rawInput := prob.InputOptions{NUL: input.NUL}
return prob.EachInputFrom(paths, in, rawInput, func(record []byte) error {
item, err := fields.Select(record)
if err != nil {
return err
}
if input.Trim {
item = bytes.TrimSpace(item)
}
if input.IgnoreEmpty && len(item) == 0 {
return nil
}
output := item
if fields.Enabled() {
output = record
}
return fn(item, output)
})
}

func writeItem(out *bufio.Writer, item []byte, nul bool) error {
if _, err := out.Write(item); err != nil {
return err
Expand Down
6 changes: 6 additions & 0 deletions cmd/heavy/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ func main() {
fmt.Fprintf(flag.CommandLine.Output(), `Usage: %s [options] [file...]

Find frequent values in newline-delimited input.
Use --delimiter and --field to rank one field per record.
Default mode uses bounded-memory approximate counts. Use --exact only when all distinct values fit in memory.

Examples:
Expand Down Expand Up @@ -64,6 +65,11 @@ Options:
os.Exit(2)
}

if err := input.Fields.Validate(); err != nil {
fmt.Fprintf(os.Stderr, "heavy: %v\n", err)
os.Exit(2)
}

results, err := heavy.Run(flag.Args(), heavy.Config{
Top: *top,
Capacity: *capacity,
Expand Down
34 changes: 21 additions & 13 deletions cmd/hll/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"encoding/json"
"flag"
"fmt"
"io"
"os"
"strings"

Expand All @@ -28,9 +29,9 @@ func main() {
var err error
switch os.Args[1] {
case "count":
err = count(os.Args[2:])
err = count(os.Args[2:], os.Stdin, os.Stdout)
case "build":
err = build(os.Args[2:])
err = build(os.Args[2:], os.Stdin, os.Stdout)
case "estimate":
err = estimate(os.Args[2:])
case "merge":
Expand Down Expand Up @@ -75,7 +76,7 @@ Run "hll <command> -h" for command-specific flags.
`)
}

func count(args []string) error {
func count(args []string, in io.Reader, out io.Writer) error {
fs := flag.NewFlagSet("hll count", flag.ExitOnError)
precision := fs.Uint("precision", uint(hll.DefaultP), "HLL precision, 4..20")
jsonOut := fs.Bool("json", false, "write JSON output")
Expand All @@ -85,7 +86,8 @@ func count(args []string) error {
fmt.Fprint(fs.Output(), `Usage: hll count [options] [file...]

Read newline-delimited values from files or stdin and print an approximate unique count.
Empty lines and surrounding whitespace are significant unless input flags change that.
Use --delimiter and --field to count one field per record.
Empty values and surrounding whitespace are significant unless input flags change that.

Example:
awk '{print $1}' access.log | hll count --ignore-empty
Expand All @@ -98,6 +100,9 @@ Options:
return err
}

if err := input.Fields.Validate(); err != nil {
return fmt.Errorf("usage: %w", err)
}
p, err := hll.Precision(*precision)
if err != nil {
return err
Expand All @@ -106,16 +111,16 @@ Options:
if err != nil {
return err
}
if err := prob.EachInput(fs.Args(), input, func(item []byte) error {
if err := prob.EachInputFrom(fs.Args(), in, input, func(item []byte) error {
s.Add(item)
return nil
}); err != nil {
return err
}
return writeEstimate(os.Stdout, s, *jsonOut)
return writeEstimate(out, s, *jsonOut)
}

func build(args []string) error {
func build(args []string, in io.Reader, out io.Writer) error {
fs := flag.NewFlagSet("hll build", flag.ExitOnError)
precision := fs.Uint("precision", uint(hll.DefaultP), "HLL precision, 4..20")
var input prob.InputOptions
Expand All @@ -124,6 +129,7 @@ func build(args []string) error {
fmt.Fprint(fs.Output(), `Usage: hll build [options] [file...] > file.hll

Read values from files or stdin and write a binary HyperLogLog sketch to stdout.
Use --delimiter and --field to hash one field per record.
Redirect stdout to save the sketch.

Options:
Expand All @@ -134,6 +140,9 @@ Options:
return err
}

if err := input.Fields.Validate(); err != nil {
return fmt.Errorf("usage: %w", err)
}
p, err := hll.Precision(*precision)
if err != nil {
return err
Expand All @@ -142,13 +151,13 @@ Options:
if err != nil {
return err
}
if err := prob.EachInput(fs.Args(), input, func(item []byte) error {
if err := prob.EachInputFrom(fs.Args(), in, input, func(item []byte) error {
s.Add(item)
return nil
}); err != nil {
return err
}
return hll.Write(os.Stdout, s)
return hll.Write(out, s)
}

func estimate(args []string) error {
Expand Down Expand Up @@ -271,15 +280,14 @@ func countStdin(paths []string) int {
return count
}

func writeEstimate(out *os.File, s *hll.Sketch, jsonOut bool) error {
func writeEstimate(out io.Writer, s *hll.Sketch, jsonOut bool) error {
m := s.Metadata()
if jsonOut {
return json.NewEncoder(out).Encode(struct {
ApproxUnique uint64 `json:"approx_unique"`
RelativeError float64 `json:"relative_error"`
}{m.ApproxUnique, m.RelativeError})
}
fmt.Fprintf(out, "approx_unique=%d\n", m.ApproxUnique)
fmt.Fprintf(out, "relative_error=%.2f%%\n", m.RelativeError*100)
return nil
_, err := fmt.Fprintf(out, "approx_unique=%d\nrelative_error=%.2f%%\n", m.ApproxUnique, m.RelativeError*100)
return err
}
68 changes: 68 additions & 0 deletions cmd/hll/main_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
package main

import (
"bytes"
"errors"
"io"
"strings"
"testing"

"github.com/jo-cube/toolbox/internal/hll"
)

func TestCountAndBuildSelectTheSameValues(t *testing.T) {
t.Parallel()
for _, nul := range []bool{false, true} {
args := []string{"--delimiter", "::", "--field", "2", "--trim", "--ignore-empty"}
input := "1:: a ::x\n2::b::y\n3::a::z\n4:: ::empty\n5::c"
if nul {
args = append(args, "-0")
input = strings.ReplaceAll(input, "\n", "\x00")
}
var direct, artifact, saved bytes.Buffer
if err := count(args, strings.NewReader(input), &direct); err != nil {
t.Fatal(err)
}
if err := build(args, strings.NewReader(input), &artifact); err != nil {
t.Fatal(err)
}
sketch, err := hll.Read(&artifact)
if err != nil {
t.Fatal(err)
}
if err := writeEstimate(&saved, sketch, false); err != nil {
t.Fatal(err)
}
if direct.String() != saved.String() || !strings.Contains(direct.String(), "approx_unique=3\n") {
t.Fatalf("direct = %q, saved = %q", direct.String(), saved.String())
}
}
}

func TestInputErrorsDoNotProduceASketchOrCount(t *testing.T) {
t.Parallel()
for _, run := range []func([]string, io.Reader, io.Writer) error{count, build} {
for _, args := range [][]string{{"--field", "2"}, {"-d", "::", "-f", "2"}} {
var out bytes.Buffer
err := run(args, strings.NewReader("missing-field\n"), &out)
if err == nil || out.Len() != 0 {
t.Fatalf("args = %q, error = %v, output bytes = %d", args, err, out.Len())
}
if len(args) == 2 && !strings.HasPrefix(err.Error(), "usage:") {
t.Fatalf("invalid field options must be a usage error: %v", err)
}
}
}
}

func TestCountPropagatesOutputFailure(t *testing.T) {
t.Parallel()
sentinel := errors.New("output failed")
if err := count(nil, strings.NewReader("a\n"), failingWriter{sentinel}); !errors.Is(err, sentinel) {
t.Fatalf("error = %v, want %v", err, sentinel)
}
}

type failingWriter struct{ err error }

func (w failingWriter) Write([]byte) (int, error) { return 0, w.err }
5 changes: 4 additions & 1 deletion cmd/sample/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ func main() {
rate := flag.Float64("rate", 0, "sample probability, 0..1")
count := flag.Int("count", 0, "reservoir sample size")
stable := flag.Bool("stable", false, "use deterministic hash sampling with --rate")
invert := flag.Bool("invert", false, "emit records excluded by the rate sample")
seed := flag.Int64("seed", 0, "random or stable hash seed")
nul := flag.Bool("nul", false, "read and write NUL-delimited records")
flag.BoolVar(nul, "0", false, "read and write NUL-delimited records")
Expand All @@ -26,7 +27,7 @@ func main() {

flag.Usage = func() {
name := filepath.Base(os.Args[0])
fmt.Fprintf(flag.CommandLine.Output(), `Usage: %s (--rate <p> [--stable] | --count <n>) [file...]
fmt.Fprintf(flag.CommandLine.Output(), `Usage: %s (--rate <p> [--stable] [--invert] | --count <n>) [file...]

Emit a subset of input records while preserving emitted records exactly.
Set exactly one of --rate or --count.
Expand All @@ -40,6 +41,7 @@ Examples:
Notes:
--rate samples each record independently unless --stable is set.
--stable hashes the full record or a selected field without its trailing delimiter.
--invert emits the complementary rate sample; use --stable or the same --seed to repeat a split.
-0 and --nul preserve NUL-delimited records instead of newline-delimited records.
--count uses reservoir sampling and writes selected records after reading input.

Expand All @@ -64,6 +66,7 @@ Options:
Count: *count,
CountSet: flagWasSet("count"),
Stable: *stable,
Invert: *invert,
Seed: *seed,
SeedSet: flagWasSet("seed"),
NUL: *nul,
Expand Down
2 changes: 1 addition & 1 deletion docs/card.md
Original file line number Diff line number Diff line change
Expand Up @@ -131,7 +131,7 @@ Separate counters report:
- `empty`: empty strings
- `total`: total records observed for that field

For non-string JSON values, `card` hashes the JSON encoding of the value. JSON numbers retain their input spelling, so `1`, `1.0`, and `1e0` are distinct and large integers keep their full precision.
For JSON values, `card` hashes the JSON encoding, preserving type distinctions: `1` and `"1"`, or `true` and `"true"`, are different values. Object key order does not affect cardinality. JSON numbers retain their input spelling, so `1`, `1.0`, and `1e0` are distinct and large integers keep their full precision.

Malformed CSV or malformed JSON fails the command with a line-aware error where possible.

Expand Down
Loading
Loading