Skip to content
Merged
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
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ require (
github.com/yamitzky/xlrd-go v0.1.0
golang.org/x/crypto v0.55.0
golang.org/x/net v0.58.0
golang.org/x/sync v0.22.0
google.golang.org/protobuf v1.36.12
gopkg.in/yaml.v3 v3.0.1
)
Expand Down Expand Up @@ -90,7 +91,6 @@ require (
github.com/xuri/efp v0.0.1 // indirect
github.com/xuri/nfp v0.0.2-0.20250530014748-2ddeb826f9a9 // indirect
golang.org/x/oauth2 v0.36.0 // indirect
golang.org/x/sync v0.22.0 // indirect
golang.org/x/sys v0.47.0 // indirect
golang.org/x/text v0.41.0 // indirect
)
Expand Down
3 changes: 2 additions & 1 deletion internal/pkg/pipeline/task/file/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@ In read mode, two values are stored in each record's context:
| `tags` | map[string]string | - | S3 **write** only: object tags applied on `PutObject`. Ignored for local paths. Values support macros and context templates. See [S3 object tags](#s3-object-tags). |
| `success_file` | bool | `false` | Whether to create a success file after writing |
| `success_file_name` | string | `_SUCCESS` | Name of the success file |
| `task_concurrency` | int | `1` | Number of competing-consumer workers for this task |
| `task_concurrency` | int | `1` | Write: competing-consumer workers. Read: parallel file fetches |
| `context` | map | - | JQ expressions whose results are stored on each record for downstream tasks |
| `fail_on_error` | bool | `false` | Whether to stop the pipeline if this task encounters an error |

Expand Down Expand Up @@ -176,6 +176,7 @@ With the inputs above, the writes land at:

- `test/pipelines/file.yaml` - Basic file operations
- `test/pipelines/context_test.yaml` - File task with context variables
- `test/pipelines/file_concurrency_test.yaml` - Concurrent read of 100 nested files (`task_concurrency: 3`)

## Use Cases

Expand Down
52 changes: 36 additions & 16 deletions internal/pkg/pipeline/task/file/file.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import (
"github.com/patterninc/caterpillar/internal/pkg/pipeline/record"
"github.com/patterninc/caterpillar/internal/pkg/pipeline/task"
"github.com/patterninc/caterpillar/internal/pkg/textutil"
"golang.org/x/sync/errgroup"
)

const (
Expand Down Expand Up @@ -131,30 +132,49 @@ func (f *file) readFile(output chan<- *record.Record) error {
return err
}

for _, path := range paths {
// if no paths are found, return nil
if len(paths) == 0 {
return nil
}

readerCloser, err := reader.read(path)
if err != nil {
return err
}
defer readerCloser.Close()
// set the error group limit to the minimum of the task concurrency and the number of paths
g, gctx := errgroup.WithContext(context.Background())
g.SetLimit(min(f.GetTaskConcurrency(), len(paths)))

content, err := io.ReadAll(readerCloser)
if err != nil {
return err
// iterate over the paths and emit the file
for _, path := range paths {
if gctx.Err() != nil {
break
}
g.Go(func() error {
return f.emitFile(reader, path, output)
})
}

// Create a default record with context
fileName := textutil.SlugifyFileName(filepath.Base(path))
rc := &record.Record{Context: ctx}
rc.SetContextValue(string(task.CtxKeyFileNameWrite), fileName)
rc.SetContextValue(string(task.CtxKeyFilePathWrite), textutil.SlugifyFilePath(path))
return g.Wait()

// let's write content to output channel
f.SendData(rc.Context, content, output)
}

func (f *file) emitFile(r reader, path string, output chan<- *record.Record) error {

readerCloser, err := r.read(path)
if err != nil {
return err
}
defer readerCloser.Close()

content, err := io.ReadAll(readerCloser)
if err != nil {
return err
}

fileName := textutil.SlugifyFileName(filepath.Base(path))
rc := &record.Record{Context: ctx}
rc.SetContextValue(string(task.CtxKeyFileNameWrite), fileName)
rc.SetContextValue(string(task.CtxKeyFilePathWrite), textutil.SlugifyFilePath(path))

f.SendData(rc.Context, content, output)

return nil

}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-0/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-0/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-0/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-0/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-0/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-1/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-1/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-1/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-1/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-1/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-2/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-2/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-2/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-2/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-2/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-3/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-3/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-3/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-3/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-0/batch-3/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-0/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-0/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-0/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-0/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-0/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-1/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-1/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-1/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-1/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-1/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-2/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-2/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-2/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-2/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-2/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-3/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-3/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-3/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-3/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-1/batch-3/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-0/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-0/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-0/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-0/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-0/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-1/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-1/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-1/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-1/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-1/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-2/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-2/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-2/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-2/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-2/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-3/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-3/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-3/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-3/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-2/batch-3/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-0/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-0/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-0/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-0/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-0/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-1/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-1/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-1/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-1/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-1/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-2/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-2/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-2/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-2/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-2/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-3/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-3/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-3/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-3/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-3/batch-3/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-0/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-0/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-0/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-0/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-0/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-1/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-1/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-1/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-1/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-1/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-2/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-2/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-2/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-2/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-2/file-4
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-3/file-0
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-3/file-1
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-3/file-2
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-3/file-3
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
region-4/batch-3/file-4
14 changes: 6 additions & 8 deletions test/pipelines/file_concurrency_test.yaml
Original file line number Diff line number Diff line change
@@ -1,11 +1,9 @@
tasks:
- name: read_glob_concurrently
- name: read_nested_concurrently
type: file
path: test/pipelines/*.txt
task_concurrency: 4
- name: split_to_lines
type: split
- name: write_concurrently
path: test/pipelines/file_concurrency/**/*.txt
task_concurrency: 3
fail_on_error: true
- name: write_to_local
type: file
path: test/pipelines/output/{{ macro "uuid" }}.txt
task_concurrency: 8
path: ./output/{{ macro "uuid" }}.txt
11 changes: 11 additions & 0 deletions test/pipelines/file_read_failure.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
# Glob matches two readable files and one directory named like a file.
# Opening that directory fails the read, and fail_on_error fails the run.
tasks:
- name: read_with_failure
type: file
path: test/pipelines/file_read_failure/*
task_concurrency: 2
fail_on_error: true
- name: echo
type: echo
only_data: true
Empty file.
1 change: 1 addition & 0 deletions test/pipelines/file_read_failure/ok-0.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
ok-0
1 change: 1 addition & 0 deletions test/pipelines/file_read_failure/ok-1.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
ok-1
Loading