From d9298fd26b16f5dbb3800aab3fd3c2d4fa88551f Mon Sep 17 00:00:00 2001 From: John Kerl Date: Wed, 26 Jan 2022 23:06:41 -0500 Subject: [PATCH] mlr split --- internal/pkg/output/file_output_handlers.go | 23 + .../pkg/transformers/aaa_transformer_table.go | 1 + internal/pkg/transformers/split.go | 437 ++++++++++++++++++ test/input/example.csv | 11 + todo.txt | 2 + 5 files changed, 474 insertions(+) create mode 100644 internal/pkg/transformers/split.go create mode 100644 test/input/example.csv diff --git a/internal/pkg/output/file_output_handlers.go b/internal/pkg/output/file_output_handlers.go index b5e1df510..cd7c3f896 100644 --- a/internal/pkg/output/file_output_handlers.go +++ b/internal/pkg/output/file_output_handlers.go @@ -56,6 +56,17 @@ type MultiOutputHandlerManager struct { } // ---------------------------------------------------------------- +func NewFileOutputHandlerManager( + recordWriterOptions *cli.TWriterOptions, + doAppend bool, +) *MultiOutputHandlerManager { + if doAppend { + return NewFileAppendHandlerManager(recordWriterOptions) + } else { + return NewFileWritetHandlerManager(recordWriterOptions) + } +} + func NewFileWritetHandlerManager( recordWriterOptions *cli.TWriterOptions, ) *MultiOutputHandlerManager { @@ -228,6 +239,18 @@ func newOutputHandlerCommon( } // ---------------------------------------------------------------- +func NewFileOutputHandler( + filename string, + recordWriterOptions *cli.TWriterOptions, + doAppend bool, +) (*FileOutputHandler, error) { + if doAppend { + return NewFileAppendOutputHandler(filename, recordWriterOptions) + } else { + return NewFileWriteOutputHandler(filename, recordWriterOptions) + } +} + func NewFileWriteOutputHandler( filename string, recordWriterOptions *cli.TWriterOptions, diff --git a/internal/pkg/transformers/aaa_transformer_table.go b/internal/pkg/transformers/aaa_transformer_table.go index ed6c0a84d..463b745a4 100644 --- a/internal/pkg/transformers/aaa_transformer_table.go +++ b/internal/pkg/transformers/aaa_transformer_table.go @@ -59,6 +59,7 @@ var TRANSFORMER_LOOKUP_TABLE = []TransformerSetup{ SkipTrivialRecordsSetup, SortSetup, SortWithinRecordsSetup, + SplitSetup, Stats1Setup, Stats2Setup, StepSetup, diff --git a/internal/pkg/transformers/split.go b/internal/pkg/transformers/split.go new file mode 100644 index 000000000..287b42768 --- /dev/null +++ b/internal/pkg/transformers/split.go @@ -0,0 +1,437 @@ +package transformers + +import ( + "bytes" + "container/list" + "fmt" + "net/url" + "os" + "strings" + + "github.com/johnkerl/miller/internal/pkg/cli" + "github.com/johnkerl/miller/internal/pkg/mlrval" + "github.com/johnkerl/miller/internal/pkg/output" + "github.com/johnkerl/miller/internal/pkg/types" +) + +// ---------------------------------------------------------------- +const verbNameSplit = "split" +const splitDefaultOutputFileNamePrefix = "split" + +var SplitSetup = TransformerSetup{ + Verb: verbNameSplit, + UsageFunc: transformerSplitUsage, + ParseCLIFunc: transformerSplitParseCLI, + IgnoresInput: false, +} + +func transformerSplitUsage( + o *os.File, + doExit bool, + exitCode int, +) { + fmt.Fprintf(o, "Usage: %s %s [options] {filename}\n", "mlr", verbNameSplit) + fmt.Fprintf(o, + `Options: +-n {n}: Cap file sizes at N records. +-m {m}: Produce M files, round-robining records among them. +-g {a,b,c}: Write separate files with records having distinct values for fields named a,b,c. +Exactly one of -m, -n, or -g must be supplied. +--prefix {p} Specify filename prefix; default "`+splitDefaultOutputFileNamePrefix+`". +--suffix {s} Specify filename suffix; default is from mlr output format, e.g. "csv". +-a Append to existing file(s), if any, rather than overwriting. +-v Send records along to downstream verbs as well as splitting to files. +-h|--help Show this message. +Any of the output-format command-line flags (see mlr -h). For example, using + mlr --icsv --from myfile.csv split --ojson -n 1000 +the input is CSV, but the output files are JSON. + +Examples: Suppose myfile.csv has 1,000,000 records. + +100 output files, 10,000 records each. First 10,000 records in split_1.csv, next in split_2.csv, etc. + mlr --csv --from myfile.csv split -n 10000 + +10 output files, 100,000 records each. Records 1,11,21,etc in split_1.csv, records 2,12,22, etc in split_2.csv, etc. + mlr --csv --from myfile.csv split -m 10 +Same, but with JSON output. + mlr --csv --from myfile.csv split -m 10 -o json + +Same but instead of split_1.csv, split_2.csv, etc. there are test_1.dat, test_2.dat, etc. + mlr --csv --from myfile.csv split -m 10 --prefix test --suffix dat +Same, but written to the /tmp/ directory. + mlr --csv --from myfile.csv split -m 10 --prefix /tmp/test --suffix dat + +If the shape field has values triangle and square, then there will be split_triangle.csv and split_square.csv. + mlr --csv --from myfile.csv split -g shape + +If the color field has values yellow and green, and the shape field has values triangle and square, +then there will be split_yellow_triangle.csv, split_yellow_square.csv, etc. + mlr --csv --from myfile.csv split -g color,shape + +See also the "tee" DSL function which lets you do more ad-hoc customization. +`) + if doExit { + os.Exit(exitCode) + } +} + +func transformerSplitParseCLI( + pargi *int, + argc int, + args []string, + mainOptions *cli.TOptions, + doConstruct bool, // false for first pass of CLI-parse, true for second pass +) IRecordTransformer { + + // Skip the verb name from the current spot in the mlr command line + argi := *pargi + verb := args[argi] + argi++ + + var n int = 0 + var doMod bool = false + var doSize bool = false + var groupByFieldNames []string = nil + var emitDownstream bool = false + var doAppend bool = false + var outputFileNamePrefix string = splitDefaultOutputFileNamePrefix + var outputFileNameSuffix string = "uninit" + haveOutputFileNameSuffix := false + + var localOptions *cli.TOptions = nil + if mainOptions != nil { + copyThereof := *mainOptions // struct copy + localOptions = ©Thereof + } + + // Parse local flags. + for argi < argc /* variable increment: 1 or 2 depending on flag */ { + opt := args[argi] + if !strings.HasPrefix(opt, "-") { + break // No more flag options to process + } + if args[argi] == "--" { + break // All transformers must do this so main-flags can follow verb-flags + } + argi++ + + if opt == "-h" || opt == "--help" { + transformerSplitUsage(os.Stdout, true, 0) + + } else if opt == "-n" { + n = cli.VerbGetIntArgOrDie(verb, opt, args, &argi, argc) + doSize = true + + } else if opt == "-m" { + n = cli.VerbGetIntArgOrDie(verb, opt, args, &argi, argc) + doMod = true + + } else if opt == "-g" { + groupByFieldNames = cli.VerbGetStringArrayArgOrDie(verb, opt, args, &argi, argc) + + } else if opt == "--prefix" { + outputFileNamePrefix = cli.VerbGetStringArgOrDie(verb, opt, args, &argi, argc) + + } else if opt == "--suffix" { + outputFileNameSuffix = cli.VerbGetStringArgOrDie(verb, opt, args, &argi, argc) + haveOutputFileNameSuffix = true + + } else if opt == "-a" { + doAppend = true + + } else if opt == "-v" { + emitDownstream = true + + } else { + // This is inelegant. For error-proofing we advance argi already in our + // loop (so individual if-statements don't need to). However, + // ParseWriterOptions expects it unadvanced. + largi := argi - 1 + if cli.FLAG_TABLE.Parse(args, argc, &largi, localOptions) { + // This lets mlr main and mlr split have different output formats. + // Nothing else to handle here. + argi = largi + } else { + transformerSplitUsage(os.Stderr, true, 1) + } + } + } + + doGroup := groupByFieldNames != nil + if !doMod && !doSize && !doGroup { + fmt.Fprintf(os.Stderr, "mlr %s: At least one of -m, -n, or -g is required.\n", verb) + os.Exit(1) + } + if (doMod && doSize) || (doMod && doGroup) || (doSize && doGroup) { + fmt.Fprintf(os.Stderr, "mlr %s: Only one of -m, -n, or -g is required.\n", verb) + os.Exit(1) + } + + cli.FinalizeWriterOptions(&localOptions.WriterOptions) + if !haveOutputFileNameSuffix { + outputFileNameSuffix = localOptions.WriterOptions.OutputFileFormat + } + + *pargi = argi + if !doConstruct { // All transformers must do this for main command-line parsing + return nil + } + + transformer, err := NewTransformerSplit( + n, + doMod, + doSize, + groupByFieldNames, + emitDownstream, + doAppend, + outputFileNamePrefix, + outputFileNameSuffix, + &localOptions.WriterOptions, + ) + if err != nil { + // Error message already printed out + os.Exit(1) + } + + return transformer +} + +// ---------------------------------------------------------------- +type TransformerSplit struct { + n int + outputFileNamePrefix string + outputFileNameSuffix string + emitDownstream bool + ungroupedCounter int + groupByFieldNames []string + recordWriterOptions *cli.TWriterOptions + doAppend bool + + // For doSize ungrouped: only one file open at a time + outputHandler output.OutputHandler + previousQuotient int + + // For all other cases: multiple files open at a time + outputHandlerManager output.OutputHandlerManager + + recordTransformerFunc RecordTransformerFunc +} + +func NewTransformerSplit( + n int, + doMod bool, + doSize bool, + groupByFieldNames []string, + emitDownstream bool, + doAppend bool, + outputFileNamePrefix string, + outputFileNameSuffix string, + recordWriterOptions *cli.TWriterOptions, +) (*TransformerSplit, error) { + + tr := &TransformerSplit{ + n: n, + outputFileNamePrefix: outputFileNamePrefix, + outputFileNameSuffix: outputFileNameSuffix, + emitDownstream: emitDownstream, + ungroupedCounter: 0, + groupByFieldNames: groupByFieldNames, + recordWriterOptions: recordWriterOptions, + doAppend: doAppend, + + outputHandler: nil, + previousQuotient: -1, + } + + tr.outputHandlerManager = output.NewFileOutputHandlerManager(recordWriterOptions, doAppend) + + if groupByFieldNames != nil { + tr.recordTransformerFunc = tr.splitGrouped + } else if doMod { + tr.recordTransformerFunc = tr.splitModUngrouped + } else { + tr.recordTransformerFunc = tr.splitSizeUngrouped + } + + return tr, nil +} + +func (tr *TransformerSplit) Transform( + inrecAndContext *types.RecordAndContext, + outputRecordsAndContexts *list.List, // list of *types.RecordAndContext + inputDownstreamDoneChannel <-chan bool, + outputDownstreamDoneChannel chan<- bool, +) { + HandleDefaultDownstreamDone(inputDownstreamDoneChannel, outputDownstreamDoneChannel) + tr.recordTransformerFunc(inrecAndContext, outputRecordsAndContexts, inputDownstreamDoneChannel, + outputDownstreamDoneChannel) +} + +func (tr *TransformerSplit) splitModUngrouped( + inrecAndContext *types.RecordAndContext, + outputRecordsAndContexts *list.List, // list of *types.RecordAndContext + inputDownstreamDoneChannel <-chan bool, + outputDownstreamDoneChannel chan<- bool, +) { + if !inrecAndContext.EndOfStream { + remainder := 1 + (tr.ungroupedCounter % tr.n) + filename := tr.makeUngroupedOutputFileName(remainder) + + err := tr.outputHandlerManager.WriteRecordAndContext(inrecAndContext, filename) + if err != nil { + fmt.Fprintf(os.Stderr, "mlr: file-write error: %v\n", err) + os.Exit(1) + } + + if tr.emitDownstream { + outputRecordsAndContexts.PushBack(inrecAndContext) + } + + tr.ungroupedCounter++ + + } else { + outputRecordsAndContexts.PushBack(inrecAndContext) // end-of-stream marker + errs := tr.outputHandlerManager.Close() + if len(errs) > 0 { + for _, err := range errs { + fmt.Fprintf(os.Stderr, "mlr: file-close error: %v\n", err) + } + os.Exit(1) + } + } +} + +func (tr *TransformerSplit) splitSizeUngrouped( + inrecAndContext *types.RecordAndContext, + outputRecordsAndContexts *list.List, // list of *types.RecordAndContext + inputDownstreamDoneChannel <-chan bool, + outputDownstreamDoneChannel chan<- bool, +) { + var err error + if !inrecAndContext.EndOfStream { + quotient := 1 + (tr.ungroupedCounter / tr.n) + + if quotient != tr.previousQuotient { + if tr.outputHandler != nil { + err = tr.outputHandler.Close() + if err != nil { + fmt.Fprintf(os.Stderr, "mlr: file-close error: %v\n", err) + os.Exit(1) + } + } + + filename := tr.makeUngroupedOutputFileName(quotient) + tr.outputHandler, err = output.NewFileOutputHandler( + filename, + tr.recordWriterOptions, + tr.doAppend, + ) + if err != nil { + fmt.Fprintf(os.Stderr, "mlr: file-open error: %v\n", err) + os.Exit(1) + } + + tr.previousQuotient = quotient + } + + err = tr.outputHandler.WriteRecordAndContext(inrecAndContext) + if err != nil { + fmt.Fprintf(os.Stderr, "mlr: file-write error: %v\n", err) + os.Exit(1) + } + + if tr.emitDownstream { + outputRecordsAndContexts.PushBack(inrecAndContext) + } + + tr.ungroupedCounter++ + + } else { + outputRecordsAndContexts.PushBack(inrecAndContext) // end-of-stream marker + + if tr.outputHandler != nil { + err := tr.outputHandler.Close() + if err != nil { + fmt.Fprintf(os.Stderr, "mlr: file-close error: %v\n", err) + os.Exit(1) + } + } + } +} + +func (tr *TransformerSplit) splitGrouped( + inrecAndContext *types.RecordAndContext, + outputRecordsAndContexts *list.List, // list of *types.RecordAndContext + inputDownstreamDoneChannel <-chan bool, + outputDownstreamDoneChannel chan<- bool, +) { + if !inrecAndContext.EndOfStream { + var filename string + groupByFieldValues, ok := inrecAndContext.Record.GetSelectedValues(tr.groupByFieldNames) + if !ok { + filename = fmt.Sprintf("%s_ungrouped.%s", tr.outputFileNamePrefix, tr.outputFileNameSuffix) + } else { + filename = tr.makeGroupedOutputFileName(groupByFieldValues) + } + err := tr.outputHandlerManager.WriteRecordAndContext(inrecAndContext, filename) + if err != nil { + fmt.Fprintf(os.Stderr, "mlr: %v\n", err) + os.Exit(1) + } + + if tr.emitDownstream { + outputRecordsAndContexts.PushBack(inrecAndContext) + } + + } else { + outputRecordsAndContexts.PushBack(inrecAndContext) // emit end-of-stream marker + + errs := tr.outputHandlerManager.Close() + if len(errs) > 0 { + for _, err := range errs { + fmt.Fprintf(os.Stderr, "mlr: file-close error: %v\n", err) + } + os.Exit(1) + } + } +} + +// makeUngroupedOutputFileName example: "split_53.csv" +func (tr *TransformerSplit) makeUngroupedOutputFileName(k int) string { + return fmt.Sprintf("%s_%d.%s", tr.outputFileNamePrefix, k, tr.outputFileNameSuffix) +} + +// makeGroupedOutputFileName example: "split_orange.csv" +func (tr *TransformerSplit) makeGroupedOutputFileName( + groupByFieldValues []*mlrval.Mlrval, +) string { + var buffer bytes.Buffer + buffer.WriteString(tr.outputFileNamePrefix) + for _, groupByFieldValue := range groupByFieldValues { + buffer.WriteString("_") + buffer.WriteString(url.QueryEscape(groupByFieldValue.String())) + } + buffer.WriteString(".") + buffer.WriteString(tr.outputFileNameSuffix) + return buffer.String() +} + +// makeGroupedIndexedOutputFileName example: "split_yellow_53.csv" +func (tr *TransformerSplit) makeGroupedIndexedOutputFileName( + groupByFieldValues []*mlrval.Mlrval, + index int, +) string { + // URL-escape the fields which come from data and which may have '/' + // etc within. Don't URL-escape the prefix since people may want to + // use prefixes like '/tmp/split' to write to the /tmp directory, etc. + var buffer bytes.Buffer + buffer.WriteString(tr.outputFileNamePrefix) + for _, groupByFieldValue := range groupByFieldValues { + buffer.WriteString("_") + buffer.WriteString(url.QueryEscape(groupByFieldValue.String())) + } + buffer.WriteString(fmt.Sprintf("_%d", index)) + buffer.WriteString(".") + buffer.WriteString(tr.outputFileNameSuffix) + return buffer.String() +} diff --git a/test/input/example.csv b/test/input/example.csv new file mode 100644 index 000000000..bf79dd5f7 --- /dev/null +++ b/test/input/example.csv @@ -0,0 +1,11 @@ +color,shape,flag,k,index,quantity,rate +yellow,triangle,true,1,11,43.6498,9.8870 +red,square,true,2,15,79.2778,0.0130 +red,circle,true,3,16,13.8103,2.9010 +red,square,false,4,48,77.5542,7.4670 +purple,triangle,false,5,51,81.2290,8.5910 +red,square,false,6,64,77.1991,9.5310 +purple,triangle,false,7,65,80.1405,5.8240 +yellow,circle,true,8,73,63.9785,4.2370 +yellow,circle,true,9,87,63.5058,8.3350 +purple,square,false,10,91,72.3735,8.2430 diff --git a/todo.txt b/todo.txt index 92fda4e5f..a6e7a7439 100644 --- a/todo.txt +++ b/todo.txt @@ -26,6 +26,7 @@ FEATURES o format/unformat o strmatch o =~ +* separate examples from FAQs ---------------------------------------------------------------- k better print-interpolate with {} etc @@ -42,6 +43,7 @@ mlr split ... -n, -g -- ? ---------------------------------------------------------------- * new example entry, with ccump and pgr + o slwin --prune (or somesuch) to only emit averages over full windows -- ? * make a lag-by-n and lead-by-n ----------------------------------------------------------------