-
Notifications
You must be signed in to change notification settings - Fork 4.7k
Expand file tree
/
Copy pathmatch.go
More file actions
419 lines (350 loc) · 11.2 KB
/
Copy pathmatch.go
File metadata and controls
419 lines (350 loc) · 11.2 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
// Licensed to the Apache Software Foundation (ASF) under one or more
// contributor license agreements. See the NOTICE file distributed with
// this work for additional information regarding copyright ownership.
// The ASF licenses this file to You under the Apache License, Version 2.0
// (the "License"); you may not use this file except in compliance with
// the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
// Package fileio provides transforms for matching and reading files.
package fileio
import (
"context"
"fmt"
"strings"
"time"
"github.com/apache/beam/sdks/v2/go/pkg/beam"
"github.com/apache/beam/sdks/v2/go/pkg/beam/core/graph/mtime"
"github.com/apache/beam/sdks/v2/go/pkg/beam/core/graph/window"
"github.com/apache/beam/sdks/v2/go/pkg/beam/core/state"
"github.com/apache/beam/sdks/v2/go/pkg/beam/io/filesystem"
"github.com/apache/beam/sdks/v2/go/pkg/beam/log"
"github.com/apache/beam/sdks/v2/go/pkg/beam/register"
"github.com/apache/beam/sdks/v2/go/pkg/beam/transforms/periodic"
)
func init() {
register.DoFn3x1[context.Context, string, func(FileMetadata), error](&matchFn{})
register.DoFn2x0[[]byte, func(string)](&matchContFn{})
register.DoFn4x1[state.Provider, string, FileMetadata, func(FileMetadata), error](
&dedupFn{},
)
register.DoFn4x1[state.Provider, string, FileMetadata, func(FileMetadata), error](
&dedupUnmodifiedFn{},
)
register.Emitter1[FileMetadata]()
register.Emitter1[string]()
register.Function1x2[FileMetadata, string, FileMetadata](keyByPath)
}
// emptyTreatment controls how empty matches of a pattern are treated.
type emptyTreatment int
const (
// emptyAllow allows empty matches.
emptyAllow emptyTreatment = iota
// emptyDisallow disallows empty matches.
emptyDisallow
// emptyAllowIfWildcard allows empty matches if the pattern contains a wildcard.
emptyAllowIfWildcard
)
type matchOption struct {
EmptyTreatment emptyTreatment
}
// MatchOptionFn is a function that can be passed to MatchFiles or MatchAll to configure options for
// matching files.
type MatchOptionFn func(*matchOption)
// MatchEmptyAllowIfWildcard specifies that empty matches are allowed if the pattern contains a
// wildcard.
func MatchEmptyAllowIfWildcard() MatchOptionFn {
return func(o *matchOption) {
o.EmptyTreatment = emptyAllowIfWildcard
}
}
// MatchEmptyAllow specifies that empty matches are allowed.
func MatchEmptyAllow() MatchOptionFn {
return func(o *matchOption) {
o.EmptyTreatment = emptyAllow
}
}
// MatchEmptyDisallow specifies that empty matches are not allowed.
func MatchEmptyDisallow() MatchOptionFn {
return func(o *matchOption) {
o.EmptyTreatment = emptyDisallow
}
}
// MatchFiles finds all files matching the glob pattern and returns a PCollection<FileMetadata> of
// the matching files. MatchFiles accepts a variadic number of MatchOptionFn that can be used to
// configure the treatment of empty matches. By default, empty matches are allowed if the pattern
// contains a wildcard.
func MatchFiles(s beam.Scope, glob string, opts ...MatchOptionFn) beam.PCollection {
s = s.Scope("fileio.MatchFiles")
filesystem.ValidateScheme(glob)
return MatchAll(s, beam.Create(s, glob), opts...)
}
// MatchAll finds all files matching the glob patterns given by the incoming PCollection<string> and
// returns a PCollection<FileMetadata> of the matching files. MatchAll accepts a variadic number of
// MatchOptionFn that can be used to configure the treatment of empty matches. By default, empty
// matches are allowed if the pattern contains a wildcard.
func MatchAll(s beam.Scope, col beam.PCollection, opts ...MatchOptionFn) beam.PCollection {
s = s.Scope("fileio.MatchAll")
option := &matchOption{
EmptyTreatment: emptyAllowIfWildcard,
}
for _, opt := range opts {
opt(option)
}
return beam.ParDo(s, newMatchFn(option), col)
}
type matchFn struct {
EmptyTreatment emptyTreatment
}
func newMatchFn(option *matchOption) *matchFn {
return &matchFn{
EmptyTreatment: option.EmptyTreatment,
}
}
func (fn *matchFn) ProcessElement(
ctx context.Context,
glob string,
emit func(FileMetadata),
) error {
if strings.TrimSpace(glob) == "" {
return nil
}
fs, err := filesystem.New(ctx, glob)
if err != nil {
return err
}
defer fs.Close()
files, err := fs.List(ctx, glob)
if err != nil {
return err
}
if len(files) == 0 {
if !allowEmptyMatch(glob, fn.EmptyTreatment) {
return fmt.Errorf("no files matching pattern %q", glob)
}
return nil
}
metadata, err := metadataFromFiles(ctx, fs, files)
if err != nil {
return err
}
for _, md := range metadata {
emit(md)
}
return nil
}
func allowEmptyMatch(glob string, treatment emptyTreatment) bool {
switch treatment {
case emptyDisallow:
return false
case emptyAllowIfWildcard:
return strings.Contains(glob, "*")
default:
return true
}
}
func metadataFromFiles(
ctx context.Context,
fs filesystem.Interface,
files []string,
) ([]FileMetadata, error) {
if len(files) == 0 {
return nil, nil
}
metadata := make([]FileMetadata, len(files))
for i, path := range files {
size, err := fs.Size(ctx, path)
if err != nil {
return nil, err
}
mTime, err := lastModified(ctx, fs, path)
if err != nil {
return nil, err
}
metadata[i] = FileMetadata{
Path: path,
Size: size,
LastModified: mTime,
}
}
return metadata, nil
}
func lastModified(ctx context.Context, fs filesystem.Interface, path string) (time.Time, error) {
lmGetter, ok := fs.(filesystem.LastModifiedGetter)
if !ok {
log.Warnf(ctx, "Filesystem %T does not implement filesystem.LastModifiedGetter", fs)
return time.Time{}, nil
}
mTime, err := lmGetter.LastModified(ctx, path)
if err != nil {
return time.Time{}, fmt.Errorf("error getting last modified time for %q: %v", path, err)
}
return mTime, nil
}
// duplicateTreatment controls how duplicate matches are treated.
type duplicateTreatment int
const (
// duplicateAllow allows duplicate matches.
duplicateAllow duplicateTreatment = iota
// duplicateAllowIfModified allows duplicate matches only if the file has been modified since it
// was last observed.
duplicateAllowIfModified
// duplicateSkip skips duplicate matches.
duplicateSkip
)
type matchContOption struct {
Start time.Time
End time.Time
DuplicateTreatment duplicateTreatment
ApplyWindow bool
}
// MatchContOptionFn is a function that can be passed to MatchContinuously to configure options for
// matching files.
type MatchContOptionFn func(*matchContOption)
// MatchStart specifies the start time for matching files.
func MatchStart(start time.Time) MatchContOptionFn {
return func(o *matchContOption) {
o.Start = start
}
}
// MatchEnd specifies the end time for matching files.
func MatchEnd(end time.Time) MatchContOptionFn {
return func(o *matchContOption) {
o.End = end
}
}
// MatchDuplicateAllow specifies that file path matches will not be deduplicated.
func MatchDuplicateAllow() MatchContOptionFn {
return func(o *matchContOption) {
o.DuplicateTreatment = duplicateAllow
}
}
// MatchDuplicateAllowIfModified specifies that file path matches will be deduplicated unless the
// file has been modified since it was last observed.
func MatchDuplicateAllowIfModified() MatchContOptionFn {
return func(o *matchContOption) {
o.DuplicateTreatment = duplicateAllowIfModified
}
}
// MatchDuplicateSkip specifies that file path matches will be deduplicated.
func MatchDuplicateSkip() MatchContOptionFn {
return func(o *matchContOption) {
o.DuplicateTreatment = duplicateSkip
}
}
// MatchApplyWindow specifies that each element will be assigned to an individual window.
func MatchApplyWindow() MatchContOptionFn {
return func(o *matchContOption) {
o.ApplyWindow = true
}
}
// MatchContinuously finds all files matching the glob pattern at the given interval and returns a
// PCollection<FileMetadata> of the matching files. MatchContinuously accepts a variadic number of
// MatchContOptionFn that can be used to configure:
//
// - Start: start time for matching files. Defaults to the current timestamp
// - End: end time for matching files. Defaults to the maximum timestamp
// - DuplicateAllow: allow emitting matches that have already been observed. Defaults to false
// - DuplicateAllowIfModified: allow emitting matches that have already been observed if the file
// has been modified since the last observation. Defaults to false
// - DuplicateSkip: skip emitting matches that have already been observed. Defaults to true
// - ApplyWindow: assign each element to an individual window with a fixed size equivalent to the
// interval. Defaults to false, i.e. all elements will reside in the global window
func MatchContinuously(
s beam.Scope,
glob string,
interval time.Duration,
opts ...MatchContOptionFn,
) beam.PCollection {
s = s.Scope("fileio.MatchContinuously")
filesystem.ValidateScheme(glob)
option := &matchContOption{
Start: mtime.Now().ToTime(),
End: mtime.MaxTimestamp.ToTime(),
ApplyWindow: false,
DuplicateTreatment: duplicateSkip,
}
for _, opt := range opts {
opt(option)
}
imp := periodic.Impulse(s, option.Start, option.End, interval, false)
globs := beam.ParDo(s, &matchContFn{Glob: glob}, imp)
matches := MatchAll(s, globs, MatchEmptyAllow())
out := dedupIfRequired(s, matches, option.DuplicateTreatment)
if option.ApplyWindow {
return beam.WindowInto(s, window.NewFixedWindows(interval), out)
}
return out
}
func dedupIfRequired(
s beam.Scope,
col beam.PCollection,
treatment duplicateTreatment,
) beam.PCollection {
if treatment == duplicateAllow {
return col
}
keyed := beam.ParDo(s, keyByPath, col)
if treatment == duplicateAllowIfModified {
return beam.ParDo(s, &dedupUnmodifiedFn{}, keyed)
}
return beam.ParDo(s, &dedupFn{}, keyed)
}
type matchContFn struct {
Glob string
}
func (fn *matchContFn) ProcessElement(_ []byte, emit func(string)) {
emit(fn.Glob)
}
func keyByPath(md FileMetadata) (string, FileMetadata) {
return md.Path, md
}
type dedupFn struct {
State state.Value[struct{}]
}
func (fn *dedupFn) ProcessElement(
sp state.Provider,
_ string,
md FileMetadata,
emit func(FileMetadata),
) error {
_, ok, err := fn.State.Read(sp)
if err != nil {
return fmt.Errorf("error reading state: %v", err)
}
if !ok {
emit(md)
if err := fn.State.Write(sp, struct{}{}); err != nil {
return fmt.Errorf("error writing state: %v", err)
}
}
return nil
}
type dedupUnmodifiedFn struct {
State state.Value[int64]
}
func (fn *dedupUnmodifiedFn) ProcessElement(
sp state.Provider,
_ string,
md FileMetadata,
emit func(FileMetadata),
) error {
prevMTime, ok, err := fn.State.Read(sp)
if err != nil {
return fmt.Errorf("error reading state: %v", err)
}
mTime := md.LastModified.UnixMilli()
if !ok || mTime > prevMTime {
emit(md)
if err := fn.State.Write(sp, mTime); err != nil {
return fmt.Errorf("error writing state: %v", err)
}
}
return nil
}