See
<https://builds.apache.org/job/beam_PostCommit_Python2/1974/display/redirect?page=changes>
Changes:
[robertwb] [BEAM-9433] Create expansion service artifact for common Java IOs.
[robertwb] Log in a daemon thread.
------------------------------------------
[...truncated 6.22 MB...]
DEBUG:apache_beam.runners.portability.fn_api_runner_transforms:Stages:
['ref_AppliedPTransform_read/Read/_SDFBoundedSourceWrapper/Impulse_5\n
read/Read/_SDFBoundedSourceWrapper/Impulse:beam:transform:impulse:v1\n must
follow: \n downstream_side_inputs: ref_PCollection_PCollection_4',
'read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction\n
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction:beam:transform:sdf_pair_with_restriction:v1\n
must follow: \n downstream_side_inputs: ref_PCollection_PCollection_4',
'read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction\n
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction:beam:transform:sdf_split_and_size_restrictions:v1\n
must follow: \n downstream_side_inputs: ref_PCollection_PCollection_4',
'read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process\n
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process:beam:transform:sdf_process_sized_element_and_restrictions:v1\n
must follow:
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction\n
downstream_side_inputs: ref_PCollection_PCollection_4',
'ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9\n
read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough):beam:transform:pardo:v1\n
must follow: \n downstream_side_inputs: ref_PCollection_PCollection_4',
'ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Impulse_11\n
read/_PassThroughThenCleanup/Create/Impulse:beam:transform:impulse:v1\n must
follow: \n downstream_side_inputs: ',
'ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/FlatMap(<lambda at
core.py:2643>)_12\n read/_PassThroughThenCleanup/Create/FlatMap(<lambda at
core.py:2643>):beam:transform:pardo:v1\n must follow: \n
downstream_side_inputs: ',
'ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Map(decode)_14\n
read/_PassThroughThenCleanup/Create/Map(decode):beam:transform:pardo:v1\n must
follow: \n downstream_side_inputs: ',
'ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(RemoveJsonFiles)_15\n
read/_PassThroughThenCleanup/ParDo(RemoveJsonFiles):beam:transform:pardo:v1\n
must follow: \n downstream_side_inputs: ',
'ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18\n
write/_StreamToBigQuery/AppendDestination:beam:transform:pardo:v1\n must
follow: \n downstream_side_inputs: ',
'ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19\n
write/_StreamToBigQuery/AddInsertIds:beam:transform:pardo:v1\n must follow: \n
downstream_side_inputs: ',
'ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21\n
write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys:beam:transform:pardo:v1\n
must follow: \n downstream_side_inputs: ',
'ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23\n
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps):beam:transform:pardo:v1\n
must follow: \n downstream_side_inputs: ',
'write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write\n
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write:beam:sink:runner:0.1\n
must follow: \n downstream_side_inputs: ',
'write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read\n
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read:beam:source:runner:0.1\n
must follow:
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write\n
downstream_side_inputs: ',
'ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps)_28\n
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps):beam:transform:pardo:v1\n
must follow: \n downstream_side_inputs: ',
'ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys_29\n
write/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys:beam:transform:pardo:v1\n
must follow: \n downstream_side_inputs: ',
'ref_AppliedPTransform_write/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn)_31\n
write/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn):beam:transform:pardo:v1\n
must follow: \n downstream_side_inputs: ']
INFO:apache_beam.runners.portability.fn_api_runner_transforms:====================
<function sink_flattens at 0x7f3404790758> ====================
DEBUG:apache_beam.runners.portability.fn_api_runner_transforms:18 [1, 1, 1, 1,
1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1]
DEBUG:apache_beam.runners.portability.fn_api_runner_transforms:Stages:
['ref_AppliedPTransform_read/Read/_SDFBoundedSourceWrapper/Impulse_5\n
read/Read/_SDFBoundedSourceWrapper/Impulse:beam:transform:impulse:v1\n must
follow: \n downstream_side_inputs: ref_PCollection_PCollection_4',
'read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction\n
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction:beam:transform:sdf_pair_with_restriction:v1\n
must follow: \n downstream_side_inputs: ref_PCollection_PCollection_4',
'read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction\n
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction:beam:transform:sdf_split_and_size_restrictions:v1\n
must follow: \n downstream_side_inputs: ref_PCollection_PCollection_4',
'read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process\n
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process:beam:transform:sdf_process_sized_element_and_restrictions:v1\n
must follow:
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction\n
downstream_side_inputs: ref_PCollection_PCollection_4',
'ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9\n
read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough):beam:transform:pardo:v1\n
must follow: \n downstream_side_inputs: ref_PCollection_PCollection_4',
'ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Impulse_11\n
read/_PassThroughThenCleanup/Create/Impulse:beam:transform:impulse:v1\n must
follow: \n downstream_side_inputs: ',
'ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/FlatMap(<lambda at
core.py:2643>)_12\n read/_PassThroughThenCleanup/Create/FlatMap(<lambda at
core.py:2643>):beam:transform:pardo:v1\n must follow: \n
downstream_side_inputs: ',
'ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Map(decode)_14\n
read/_PassThroughThenCleanup/Create/Map(decode):beam:transform:pardo:v1\n must
follow: \n downstream_side_inputs: ',
'ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(RemoveJsonFiles)_15\n
read/_PassThroughThenCleanup/ParDo(RemoveJsonFiles):beam:transform:pardo:v1\n
must follow: \n downstream_side_inputs: ',
'ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18\n
write/_StreamToBigQuery/AppendDestination:beam:transform:pardo:v1\n must
follow: \n downstream_side_inputs: ',
'ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19\n
write/_StreamToBigQuery/AddInsertIds:beam:transform:pardo:v1\n must follow: \n
downstream_side_inputs: ',
'ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21\n
write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys:beam:transform:pardo:v1\n
must follow: \n downstream_side_inputs: ',
'ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23\n
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps):beam:transform:pardo:v1\n
must follow: \n downstream_side_inputs: ',
'write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write\n
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write:beam:sink:runner:0.1\n
must follow: \n downstream_side_inputs: ',
'write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read\n
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read:beam:source:runner:0.1\n
must follow:
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write\n
downstream_side_inputs: ',
'ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps)_28\n
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps):beam:transform:pardo:v1\n
must follow: \n downstream_side_inputs: ',
'ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys_29\n
write/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys:beam:transform:pardo:v1\n
must follow: \n downstream_side_inputs: ',
'ref_AppliedPTransform_write/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn)_31\n
write/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn):beam:transform:pardo:v1\n
must follow: \n downstream_side_inputs: ']
INFO:apache_beam.runners.portability.fn_api_runner_transforms:====================
<function greedily_fuse at 0x7f34047907d0> ====================
DEBUG:apache_beam.runners.portability.fn_api_runner_transforms:4 [9, 4, 4, 4]
DEBUG:apache_beam.runners.portability.fn_api_runner_transforms:Stages:
['(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
ref_PCollection_PCollection_1_split/Read:beam:source:runner:0.1\nread/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process:beam:transform:sdf_process_sized_element_and_restrictions:v1\nread/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough):beam:transform:pardo:v1\nref_PCollection_PCollection_4/Write:beam:sink:runner:0.1\nwrite/_StreamToBigQuery/AppendDestination:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/AddInsertIds:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/AddRandomKeys:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps):beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write:beam:sink:runner:0.1\n
must follow:
((ref_AppliedPTransform_read/Read/_SDFBoundedSourceWrapper/Impulse_5)+(read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction))+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction)+(ref_PCollection_PCollection_1_split/Write))\n
downstream_side_inputs: ref_PCollection_PCollection_4',
'(ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Impulse_11)+((ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/FlatMap(<lambda
at
core.py:2643>)_12)+((ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Map(decode)_14)+(ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(RemoveJsonFiles)_15)))\n
read/_PassThroughThenCleanup/Create/Impulse:beam:transform:impulse:v1\nread/_PassThroughThenCleanup/Create/FlatMap(<lambda
at
core.py:2643>):beam:transform:pardo:v1\nread/_PassThroughThenCleanup/Create/Map(decode):beam:transform:pardo:v1\nread/_PassThroughThenCleanup/ParDo(RemoveJsonFiles):beam:transform:pardo:v1\n
must follow:
(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
downstream_side_inputs: ',
'((ref_AppliedPTransform_read/Read/_SDFBoundedSourceWrapper/Impulse_5)+(read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction))+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction)+(ref_PCollection_PCollection_1_split/Write))\n
read/Read/_SDFBoundedSourceWrapper/Impulse:beam:transform:impulse:v1\nread/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction:beam:transform:sdf_pair_with_restriction:v1\nread/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction:beam:transform:sdf_split_and_size_restrictions:v1\nref_PCollection_PCollection_1_split/Write:beam:sink:runner:0.1\n
must follow: \n downstream_side_inputs: ref_PCollection_PCollection_4',
'((write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read)+(ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps)_28))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys_29)+(ref_AppliedPTransform_write/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn)_31))\n
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read:beam:source:runner:0.1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps):beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn):beam:transform:pardo:v1\n
must follow:
(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
downstream_side_inputs: ']
INFO:apache_beam.runners.portability.fn_api_runner_transforms:====================
<function read_to_impulse at 0x7f3404790848> ====================
DEBUG:apache_beam.runners.portability.fn_api_runner_transforms:4 [9, 4, 4, 4]
DEBUG:apache_beam.runners.portability.fn_api_runner_transforms:Stages:
['(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
ref_PCollection_PCollection_1_split/Read:beam:source:runner:0.1\nread/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process:beam:transform:sdf_process_sized_element_and_restrictions:v1\nread/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough):beam:transform:pardo:v1\nref_PCollection_PCollection_4/Write:beam:sink:runner:0.1\nwrite/_StreamToBigQuery/AppendDestination:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/AddInsertIds:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/AddRandomKeys:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps):beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write:beam:sink:runner:0.1\n
must follow:
((ref_AppliedPTransform_read/Read/_SDFBoundedSourceWrapper/Impulse_5)+(read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction))+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction)+(ref_PCollection_PCollection_1_split/Write))\n
downstream_side_inputs: ref_PCollection_PCollection_4',
'(ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Impulse_11)+((ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/FlatMap(<lambda
at
core.py:2643>)_12)+((ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Map(decode)_14)+(ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(RemoveJsonFiles)_15)))\n
read/_PassThroughThenCleanup/Create/Impulse:beam:transform:impulse:v1\nread/_PassThroughThenCleanup/Create/FlatMap(<lambda
at
core.py:2643>):beam:transform:pardo:v1\nread/_PassThroughThenCleanup/Create/Map(decode):beam:transform:pardo:v1\nread/_PassThroughThenCleanup/ParDo(RemoveJsonFiles):beam:transform:pardo:v1\n
must follow:
(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
downstream_side_inputs: ',
'((ref_AppliedPTransform_read/Read/_SDFBoundedSourceWrapper/Impulse_5)+(read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction))+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction)+(ref_PCollection_PCollection_1_split/Write))\n
read/Read/_SDFBoundedSourceWrapper/Impulse:beam:transform:impulse:v1\nread/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction:beam:transform:sdf_pair_with_restriction:v1\nread/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction:beam:transform:sdf_split_and_size_restrictions:v1\nref_PCollection_PCollection_1_split/Write:beam:sink:runner:0.1\n
must follow: \n downstream_side_inputs: ref_PCollection_PCollection_4',
'((write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read)+(ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps)_28))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys_29)+(ref_AppliedPTransform_write/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn)_31))\n
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read:beam:source:runner:0.1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps):beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn):beam:transform:pardo:v1\n
must follow:
(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
downstream_side_inputs: ']
INFO:apache_beam.runners.portability.fn_api_runner_transforms:====================
<function impulse_to_input at 0x7f34047908c0> ====================
DEBUG:apache_beam.runners.portability.fn_api_runner_transforms:4 [9, 4, 4, 4]
DEBUG:apache_beam.runners.portability.fn_api_runner_transforms:Stages:
['(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
ref_PCollection_PCollection_1_split/Read:beam:source:runner:0.1\nread/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process:beam:transform:sdf_process_sized_element_and_restrictions:v1\nread/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough):beam:transform:pardo:v1\nref_PCollection_PCollection_4/Write:beam:sink:runner:0.1\nwrite/_StreamToBigQuery/AppendDestination:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/AddInsertIds:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/AddRandomKeys:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps):beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write:beam:sink:runner:0.1\n
must follow:
((ref_AppliedPTransform_read/Read/_SDFBoundedSourceWrapper/Impulse_5)+(read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction))+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction)+(ref_PCollection_PCollection_1_split/Write))\n
downstream_side_inputs: ref_PCollection_PCollection_4',
'(ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Impulse_11)+((ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/FlatMap(<lambda
at
core.py:2643>)_12)+((ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Map(decode)_14)+(ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(RemoveJsonFiles)_15)))\n
read/_PassThroughThenCleanup/Create/FlatMap(<lambda at
core.py:2643>):beam:transform:pardo:v1\nread/_PassThroughThenCleanup/Create/Map(decode):beam:transform:pardo:v1\nread/_PassThroughThenCleanup/ParDo(RemoveJsonFiles):beam:transform:pardo:v1\nread/_PassThroughThenCleanup/Create/Impulse:beam:source:runner:0.1\n
must follow:
(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
downstream_side_inputs: ',
'((ref_AppliedPTransform_read/Read/_SDFBoundedSourceWrapper/Impulse_5)+(read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction))+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction)+(ref_PCollection_PCollection_1_split/Write))\n
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction:beam:transform:sdf_pair_with_restriction:v1\nread/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction:beam:transform:sdf_split_and_size_restrictions:v1\nref_PCollection_PCollection_1_split/Write:beam:sink:runner:0.1\nread/Read/_SDFBoundedSourceWrapper/Impulse:beam:source:runner:0.1\n
must follow: \n downstream_side_inputs: ref_PCollection_PCollection_4',
'((write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read)+(ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps)_28))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys_29)+(ref_AppliedPTransform_write/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn)_31))\n
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read:beam:source:runner:0.1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps):beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn):beam:transform:pardo:v1\n
must follow:
(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
downstream_side_inputs: ']
INFO:apache_beam.runners.portability.fn_api_runner_transforms:====================
<function inject_timer_pcollections at 0x7f3404790a28> ====================
DEBUG:apache_beam.runners.portability.fn_api_runner_transforms:4 [9, 4, 4, 4]
DEBUG:apache_beam.runners.portability.fn_api_runner_transforms:Stages:
['(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
ref_PCollection_PCollection_1_split/Read:beam:source:runner:0.1\nread/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process:beam:transform:sdf_process_sized_element_and_restrictions:v1\nread/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough):beam:transform:pardo:v1\nref_PCollection_PCollection_4/Write:beam:sink:runner:0.1\nwrite/_StreamToBigQuery/AppendDestination:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/AddInsertIds:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/AddRandomKeys:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps):beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write:beam:sink:runner:0.1\n
must follow:
((ref_AppliedPTransform_read/Read/_SDFBoundedSourceWrapper/Impulse_5)+(read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction))+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction)+(ref_PCollection_PCollection_1_split/Write))\n
downstream_side_inputs: ref_PCollection_PCollection_4',
'(ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Impulse_11)+((ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/FlatMap(<lambda
at
core.py:2643>)_12)+((ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Map(decode)_14)+(ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(RemoveJsonFiles)_15)))\n
read/_PassThroughThenCleanup/Create/FlatMap(<lambda at
core.py:2643>):beam:transform:pardo:v1\nread/_PassThroughThenCleanup/Create/Map(decode):beam:transform:pardo:v1\nread/_PassThroughThenCleanup/ParDo(RemoveJsonFiles):beam:transform:pardo:v1\nread/_PassThroughThenCleanup/Create/Impulse:beam:source:runner:0.1\n
must follow:
(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
downstream_side_inputs: ',
'((ref_AppliedPTransform_read/Read/_SDFBoundedSourceWrapper/Impulse_5)+(read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction))+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction)+(ref_PCollection_PCollection_1_split/Write))\n
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction:beam:transform:sdf_pair_with_restriction:v1\nread/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction:beam:transform:sdf_split_and_size_restrictions:v1\nref_PCollection_PCollection_1_split/Write:beam:sink:runner:0.1\nread/Read/_SDFBoundedSourceWrapper/Impulse:beam:source:runner:0.1\n
must follow: \n downstream_side_inputs: ref_PCollection_PCollection_4',
'((write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read)+(ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps)_28))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys_29)+(ref_AppliedPTransform_write/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn)_31))\n
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read:beam:source:runner:0.1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps):beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn):beam:transform:pardo:v1\n
must follow:
(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
downstream_side_inputs: ']
INFO:apache_beam.runners.portability.fn_api_runner_transforms:====================
<function sort_stages at 0x7f3404790aa0> ====================
DEBUG:apache_beam.runners.portability.fn_api_runner_transforms:4 [4, 9, 4, 4]
DEBUG:apache_beam.runners.portability.fn_api_runner_transforms:Stages:
['((ref_AppliedPTransform_read/Read/_SDFBoundedSourceWrapper/Impulse_5)+(read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction))+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction)+(ref_PCollection_PCollection_1_split/Write))\n
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction:beam:transform:sdf_pair_with_restriction:v1\nread/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction:beam:transform:sdf_split_and_size_restrictions:v1\nref_PCollection_PCollection_1_split/Write:beam:sink:runner:0.1\nread/Read/_SDFBoundedSourceWrapper/Impulse:beam:source:runner:0.1\n
must follow: \n downstream_side_inputs: ref_PCollection_PCollection_4',
'(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
ref_PCollection_PCollection_1_split/Read:beam:source:runner:0.1\nread/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process:beam:transform:sdf_process_sized_element_and_restrictions:v1\nread/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough):beam:transform:pardo:v1\nref_PCollection_PCollection_4/Write:beam:sink:runner:0.1\nwrite/_StreamToBigQuery/AppendDestination:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/AddInsertIds:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/AddRandomKeys:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps):beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write:beam:sink:runner:0.1\n
must follow:
((ref_AppliedPTransform_read/Read/_SDFBoundedSourceWrapper/Impulse_5)+(read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction))+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction)+(ref_PCollection_PCollection_1_split/Write))\n
downstream_side_inputs: ref_PCollection_PCollection_4',
'(ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Impulse_11)+((ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/FlatMap(<lambda
at
core.py:2643>)_12)+((ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Map(decode)_14)+(ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(RemoveJsonFiles)_15)))\n
read/_PassThroughThenCleanup/Create/FlatMap(<lambda at
core.py:2643>):beam:transform:pardo:v1\nread/_PassThroughThenCleanup/Create/Map(decode):beam:transform:pardo:v1\nread/_PassThroughThenCleanup/ParDo(RemoveJsonFiles):beam:transform:pardo:v1\nread/_PassThroughThenCleanup/Create/Impulse:beam:source:runner:0.1\n
must follow:
(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
downstream_side_inputs: ',
'((write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read)+(ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps)_28))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys_29)+(ref_AppliedPTransform_write/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn)_31))\n
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read:beam:source:runner:0.1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps):beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn):beam:transform:pardo:v1\n
must follow:
(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
downstream_side_inputs: ']
INFO:apache_beam.runners.portability.fn_api_runner_transforms:====================
<function window_pcollection_coders at 0x7f3404790b18> ====================
DEBUG:apache_beam.runners.portability.fn_api_runner_transforms:4 [4, 9, 4, 4]
DEBUG:apache_beam.runners.portability.fn_api_runner_transforms:Stages:
['((ref_AppliedPTransform_read/Read/_SDFBoundedSourceWrapper/Impulse_5)+(read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction))+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction)+(ref_PCollection_PCollection_1_split/Write))\n
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction:beam:transform:sdf_pair_with_restriction:v1\nread/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction:beam:transform:sdf_split_and_size_restrictions:v1\nref_PCollection_PCollection_1_split/Write:beam:sink:runner:0.1\nread/Read/_SDFBoundedSourceWrapper/Impulse:beam:source:runner:0.1\n
must follow: \n downstream_side_inputs: ref_PCollection_PCollection_4',
'(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
ref_PCollection_PCollection_1_split/Read:beam:source:runner:0.1\nread/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process:beam:transform:sdf_process_sized_element_and_restrictions:v1\nread/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough):beam:transform:pardo:v1\nref_PCollection_PCollection_4/Write:beam:sink:runner:0.1\nwrite/_StreamToBigQuery/AppendDestination:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/AddInsertIds:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/AddRandomKeys:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps):beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write:beam:sink:runner:0.1\n
must follow:
((ref_AppliedPTransform_read/Read/_SDFBoundedSourceWrapper/Impulse_5)+(read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction))+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction)+(ref_PCollection_PCollection_1_split/Write))\n
downstream_side_inputs: ref_PCollection_PCollection_4',
'(ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Impulse_11)+((ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/FlatMap(<lambda
at
core.py:2643>)_12)+((ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Map(decode)_14)+(ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(RemoveJsonFiles)_15)))\n
read/_PassThroughThenCleanup/Create/FlatMap(<lambda at
core.py:2643>):beam:transform:pardo:v1\nread/_PassThroughThenCleanup/Create/Map(decode):beam:transform:pardo:v1\nread/_PassThroughThenCleanup/ParDo(RemoveJsonFiles):beam:transform:pardo:v1\nread/_PassThroughThenCleanup/Create/Impulse:beam:source:runner:0.1\n
must follow:
(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
downstream_side_inputs: ',
'((write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read)+(ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps)_28))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys_29)+(ref_AppliedPTransform_write/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn)_31))\n
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read:beam:source:runner:0.1\nwrite/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps):beam:transform:pardo:v1\nwrite/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys:beam:transform:pardo:v1\nwrite/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn):beam:transform:pardo:v1\n
must follow:
(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))\n
downstream_side_inputs: ']
INFO:apache_beam.runners.worker.statecache:Creating state cache with size 100
INFO:apache_beam.runners.portability.fn_api_runner:Created Worker handler
<apache_beam.runners.portability.fn_api_runner.EmbeddedWorkerHandler object at
0x7f3404158d50> for environment urn: "beam:env:embedded_python:v1"
INFO:apache_beam.runners.portability.fn_api_runner:Running
((ref_AppliedPTransform_read/Read/_SDFBoundedSourceWrapper/Impulse_5)+(read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction))+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction)+(ref_PCollection_PCollection_1_split/Write))
DEBUG:apache_beam.runners.worker.bundle_processor:start <DataOutputOperation
ref_PCollection_PCollection_1_split/Write >
DEBUG:apache_beam.runners.worker.bundle_processor:start <DoOperation
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction
output_tags=[u'out'],
receivers=[SingletonConsumerSet[read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction.out0,
coder=WindowedValueCoder[TupleCoder[TupleCoder[BytesCoder,
TupleCoder[LengthPrefixCoder[DillCoder],
LengthPrefixCoder[FastPrimitivesCoder]]], FloatCoder]], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:start <DoOperation
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction
output_tags=[u'out'],
receivers=[SingletonConsumerSet[read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction.out0,
coder=WindowedValueCoder[TupleCoder[BytesCoder, TupleCoder[DillCoder,
FastPrimitivesCoder]]], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:start <DataInputOperation
read/Read/_SDFBoundedSourceWrapper/Impulse
receivers=[SingletonConsumerSet[read/Read/_SDFBoundedSourceWrapper/Impulse.out0,
coder=WindowedValueCoder[BytesCoder], len(consumers)=1]]>
DEBUG:apache_beam.io.gcp.bigquery_tools:Query SELECT * FROM (SELECT "apple" as
fruit) UNION ALL (SELECT "orange" as fruit) does not reference any tables.
WARNING:apache_beam.io.gcp.bigquery_tools:Dataset
apache-beam-testing:temp_dataset_e49704b240a240d49840827efe8554db does not
exist so we will create it as temporary with location=None
INFO:root:Job status: RUNNING
INFO:root:Job status: DONE
INFO:root:Job status: RUNNING
INFO:root:Job status: DONE
DEBUG:apache_beam.io.filesystem:Listing files in
'gs://temp-storage-for-end-to-end-tests/temp-it/111ed58d407f48828928efb42dbc3c6b/bigquery-table-dump-'
DEBUG:apache_beam.io.filesystem:translate_pattern:
'gs://temp-storage-for-end-to-end-tests/temp-it/111ed58d407f48828928efb42dbc3c6b/bigquery-table-dump-*.json'
->
'gs\\:\\/\\/temp\\-storage\\-for\\-end\\-to\\-end\\-tests\\/temp\\-it\\/111ed58d407f48828928efb42dbc3c6b\\/bigquery\\-table\\-dump\\-[^/\\\\]*\\.json'
INFO:apache_beam.io.gcp.gcsio:Starting the size estimation of the input
INFO:apache_beam.io.gcp.gcsio:Finished listing 1 files in 0.0375981330872
seconds.
DEBUG:apache_beam.io.filesystem:translate_pattern:
'gs://temp-storage-for-end-to-end-tests/temp-it/111ed58d407f48828928efb42dbc3c6b/bigquery-table-dump-000000000000.json'
->
'gs\\:\\/\\/temp\\-storage\\-for\\-end\\-to\\-end\\-tests\\/temp\\-it\\/111ed58d407f48828928efb42dbc3c6b\\/bigquery\\-table\\-dump\\-000000000000\\.json'
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DataInputOperation
read/Read/_SDFBoundedSourceWrapper/Impulse
receivers=[SingletonConsumerSet[read/Read/_SDFBoundedSourceWrapper/Impulse.out0,
coder=WindowedValueCoder[BytesCoder], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DoOperation
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction
output_tags=[u'out'],
receivers=[SingletonConsumerSet[read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction.out0,
coder=WindowedValueCoder[TupleCoder[BytesCoder, TupleCoder[DillCoder,
FastPrimitivesCoder]]], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DoOperation
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction
output_tags=[u'out'],
receivers=[SingletonConsumerSet[read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/SplitAndSizeRestriction.out0,
coder=WindowedValueCoder[TupleCoder[TupleCoder[BytesCoder,
TupleCoder[LengthPrefixCoder[DillCoder],
LengthPrefixCoder[FastPrimitivesCoder]]], FloatCoder]], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DataOutputOperation
ref_PCollection_PCollection_1_split/Write >
DEBUG:apache_beam.runners.portability.fn_api_runner:Wait for the bundle
bundle_9 to finish.
INFO:apache_beam.runners.portability.fn_api_runner:Running
(((ref_PCollection_PCollection_1_split/Read)+((read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process)+(((ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)_9)+(ref_PCollection_PCollection_4/Write))+(ref_AppliedPTransform_write/_StreamToBigQuery/AppendDestination_18))))+(ref_AppliedPTransform_write/_StreamToBigQuery/AddInsertIds_19))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys_21)+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)_23)+(write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write)))
DEBUG:apache_beam.runners.worker.bundle_processor:start <DataOutputOperation
ref_PCollection_PCollection_4/Write >
DEBUG:apache_beam.runners.worker.bundle_processor:start <DataOutputOperation
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write >
DEBUG:apache_beam.runners.worker.bundle_processor:start <DoOperation
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)
output_tags=[u'None'],
receivers=[SingletonConsumerSet[write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps).out0,
coder=WindowedValueCoder[TupleCoder[LengthPrefixCoder[FastPrimitivesCoder],
LengthPrefixCoder[FastPrimitivesCoder]]], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:start <DoOperation
write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys output_tags=[u'None'],
receivers=[SingletonConsumerSet[write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys.out0,
coder=WindowedValueCoder[TupleCoder[FastPrimitivesCoder,
TupleCoder[FastPrimitivesCoder, TupleCoder[FastPrimitivesCoder,
FastPrimitivesCoder]]]], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:start <DoOperation
write/_StreamToBigQuery/AddInsertIds output_tags=[u'None'],
receivers=[SingletonConsumerSet[write/_StreamToBigQuery/AddInsertIds.out0,
coder=WindowedValueCoder[TupleCoder[FastPrimitivesCoder,
TupleCoder[FastPrimitivesCoder, FastPrimitivesCoder]]], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:start <DoOperation
write/_StreamToBigQuery/AppendDestination output_tags=[u'None'],
receivers=[SingletonConsumerSet[write/_StreamToBigQuery/AppendDestination.out0,
coder=WindowedValueCoder[TupleCoder[FastPrimitivesCoder, FastPrimitivesCoder]],
len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:start <DoOperation
read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)
output_tags=[u'None', u'cleanup_signal'],
receivers=[SingletonConsumerSet[read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough).out0,
coder=WindowedValueCoder[FastPrimitivesCoder], len(consumers)=1],
SingletonConsumerSet[read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough).out1,
coder=WindowedValueCoder[LengthPrefixCoder[FastPrimitivesCoder]],
len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:start
<SdfProcessSizedElements
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process
output_tags=[u'None'],
receivers=[SingletonConsumerSet[read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process.out0,
coder=WindowedValueCoder[FastPrimitivesCoder], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:start <DataInputOperation
ref_PCollection_PCollection_1_split/Read
receivers=[SingletonConsumerSet[ref_PCollection_PCollection_1_split/Read.out0,
coder=WindowedValueCoder[TupleCoder[TupleCoder[BytesCoder,
TupleCoder[LengthPrefixCoder[DillCoder],
LengthPrefixCoder[FastPrimitivesCoder]]], FloatCoder]], len(consumers)=1]]>
DEBUG:apache_beam.io.filesystem:translate_pattern:
'gs://temp-storage-for-end-to-end-tests/temp-it/111ed58d407f48828928efb42dbc3c6b/bigquery-table-dump-000000000000.json'
->
'gs\\:\\/\\/temp\\-storage\\-for\\-end\\-to\\-end\\-tests\\/temp\\-it\\/111ed58d407f48828928efb42dbc3c6b\\/bigquery\\-table\\-dump\\-000000000000\\.json'
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DataInputOperation
ref_PCollection_PCollection_1_split/Read
receivers=[SingletonConsumerSet[ref_PCollection_PCollection_1_split/Read.out0,
coder=WindowedValueCoder[TupleCoder[TupleCoder[BytesCoder,
TupleCoder[LengthPrefixCoder[DillCoder],
LengthPrefixCoder[FastPrimitivesCoder]]], FloatCoder]], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:finish
<SdfProcessSizedElements
read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process
output_tags=[u'None'],
receivers=[SingletonConsumerSet[read/Read/_SDFBoundedSourceWrapper/ParDo(SDFBoundedSourceDoFn)/Process.out0,
coder=WindowedValueCoder[FastPrimitivesCoder], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DoOperation
read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough)
output_tags=[u'None', u'cleanup_signal'],
receivers=[SingletonConsumerSet[read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough).out0,
coder=WindowedValueCoder[FastPrimitivesCoder], len(consumers)=1],
SingletonConsumerSet[read/_PassThroughThenCleanup/ParDo(PassThrough)/ParDo(PassThrough).out1,
coder=WindowedValueCoder[LengthPrefixCoder[FastPrimitivesCoder]],
len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DoOperation
write/_StreamToBigQuery/AppendDestination output_tags=[u'None'],
receivers=[SingletonConsumerSet[write/_StreamToBigQuery/AppendDestination.out0,
coder=WindowedValueCoder[TupleCoder[FastPrimitivesCoder, FastPrimitivesCoder]],
len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DoOperation
write/_StreamToBigQuery/AddInsertIds output_tags=[u'None'],
receivers=[SingletonConsumerSet[write/_StreamToBigQuery/AddInsertIds.out0,
coder=WindowedValueCoder[TupleCoder[FastPrimitivesCoder,
TupleCoder[FastPrimitivesCoder, FastPrimitivesCoder]]], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DoOperation
write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys output_tags=[u'None'],
receivers=[SingletonConsumerSet[write/_StreamToBigQuery/CommitInsertIds/AddRandomKeys.out0,
coder=WindowedValueCoder[TupleCoder[FastPrimitivesCoder,
TupleCoder[FastPrimitivesCoder, TupleCoder[FastPrimitivesCoder,
FastPrimitivesCoder]]]], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DoOperation
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps)
output_tags=[u'None'],
receivers=[SingletonConsumerSet[write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/Map(reify_timestamps).out0,
coder=WindowedValueCoder[TupleCoder[LengthPrefixCoder[FastPrimitivesCoder],
LengthPrefixCoder[FastPrimitivesCoder]]], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DataOutputOperation
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Write >
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DataOutputOperation
ref_PCollection_PCollection_4/Write >
DEBUG:apache_beam.runners.portability.fn_api_runner:Wait for the bundle
bundle_10 to finish.
INFO:apache_beam.runners.portability.fn_api_runner:Running
(ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Impulse_11)+((ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/FlatMap(<lambda
at
core.py:2643>)_12)+((ref_AppliedPTransform_read/_PassThroughThenCleanup/Create/Map(decode)_14)+(ref_AppliedPTransform_read/_PassThroughThenCleanup/ParDo(RemoveJsonFiles)_15)))
DEBUG:apache_beam.runners.worker.bundle_processor:start <DoOperation
read/_PassThroughThenCleanup/ParDo(RemoveJsonFiles) output_tags=[u'None'],
receivers=[ConsumerSet[read/_PassThroughThenCleanup/ParDo(RemoveJsonFiles).out0,
coder=WindowedValueCoder[FastPrimitivesCoder], len(consumers)=0]]>
DEBUG:apache_beam.runners.worker.bundle_processor:start <DoOperation
read/_PassThroughThenCleanup/Create/Map(decode) output_tags=[u'None'],
receivers=[SingletonConsumerSet[read/_PassThroughThenCleanup/Create/Map(decode).out0,
coder=WindowedValueCoder[FastPrimitivesCoder], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:start <DoOperation
read/_PassThroughThenCleanup/Create/FlatMap(<lambda at core.py:2643>)
output_tags=[u'None'],
receivers=[SingletonConsumerSet[read/_PassThroughThenCleanup/Create/FlatMap(<lambda
at core.py:2643>).out0, coder=WindowedValueCoder[BytesCoder],
len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:start <DataInputOperation
read/_PassThroughThenCleanup/Create/Impulse
receivers=[SingletonConsumerSet[read/_PassThroughThenCleanup/Create/Impulse.out0,
coder=WindowedValueCoder[BytesCoder], len(consumers)=1]]>
DEBUG:apache_beam.io.filesystem:Listing files in
'gs://temp-storage-for-end-to-end-tests/temp-it/111ed58d407f48828928efb42dbc3c6b/bigquery-table-dump-'
DEBUG:apache_beam.io.filesystem:translate_pattern:
'gs://temp-storage-for-end-to-end-tests/temp-it/111ed58d407f48828928efb42dbc3c6b/bigquery-table-dump-*.json'
->
'gs\\:\\/\\/temp\\-storage\\-for\\-end\\-to\\-end\\-tests\\/temp\\-it\\/111ed58d407f48828928efb42dbc3c6b\\/bigquery\\-table\\-dump\\-[^/\\\\]*\\.json'
INFO:apache_beam.io.gcp.gcsio:Starting the size estimation of the input
INFO:apache_beam.io.gcp.gcsio:Finished listing 1 files in 0.045154094696
seconds.
DEBUG:root:RemoveJsonFiles: matched 1 files
DEBUG:apache_beam.io.filesystem:translate_pattern:
u'gs://temp-storage-for-end-to-end-tests/temp-it/111ed58d407f48828928efb42dbc3c6b/bigquery-table-dump-000000000000.json'
->
u'gs\\:\\/\\/temp\\-storage\\-for\\-end\\-to\\-end\\-tests\\/temp\\-it\\/111ed58d407f48828928efb42dbc3c6b\\/bigquery\\-table\\-dump\\-000000000000\\.json'
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DataInputOperation
read/_PassThroughThenCleanup/Create/Impulse
receivers=[SingletonConsumerSet[read/_PassThroughThenCleanup/Create/Impulse.out0,
coder=WindowedValueCoder[BytesCoder], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DoOperation
read/_PassThroughThenCleanup/Create/FlatMap(<lambda at core.py:2643>)
output_tags=[u'None'],
receivers=[SingletonConsumerSet[read/_PassThroughThenCleanup/Create/FlatMap(<lambda
at core.py:2643>).out0, coder=WindowedValueCoder[BytesCoder],
len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DoOperation
read/_PassThroughThenCleanup/Create/Map(decode) output_tags=[u'None'],
receivers=[SingletonConsumerSet[read/_PassThroughThenCleanup/Create/Map(decode).out0,
coder=WindowedValueCoder[FastPrimitivesCoder], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DoOperation
read/_PassThroughThenCleanup/ParDo(RemoveJsonFiles) output_tags=[u'None'],
receivers=[ConsumerSet[read/_PassThroughThenCleanup/ParDo(RemoveJsonFiles).out0,
coder=WindowedValueCoder[FastPrimitivesCoder], len(consumers)=0]]>
DEBUG:apache_beam.runners.portability.fn_api_runner:Wait for the bundle
bundle_11 to finish.
INFO:apache_beam.runners.portability.fn_api_runner:Running
((write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read)+(ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps)_28))+((ref_AppliedPTransform_write/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys_29)+(ref_AppliedPTransform_write/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn)_31))
DEBUG:apache_beam.runners.worker.bundle_processor:start <DoOperation
write/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn)
output_tags=[u'None', u'FailedRows'],
receivers=[ConsumerSet[write/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn).out0,
coder=WindowedValueCoder[FastPrimitivesCoder], len(consumers)=0],
ConsumerSet[write/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn).out1,
coder=WindowedValueCoder[FastPrimitivesCoder], len(consumers)=0]]>
DEBUG:apache_beam.runners.worker.bundle_processor:start <DoOperation
write/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys output_tags=[u'None'],
receivers=[SingletonConsumerSet[write/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys.out0,
coder=WindowedValueCoder[TupleCoder[FastPrimitivesCoder,
TupleCoder[FastPrimitivesCoder, FastPrimitivesCoder]]], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:start <DoOperation
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps)
output_tags=[u'None'],
receivers=[SingletonConsumerSet[write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps).out0,
coder=WindowedValueCoder[TupleCoder[FastPrimitivesCoder,
TupleCoder[FastPrimitivesCoder, TupleCoder[FastPrimitivesCoder,
FastPrimitivesCoder]]]], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:start <DataInputOperation
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read
receivers=[SingletonConsumerSet[write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read.out0,
coder=WindowedValueCoder[TupleCoder[LengthPrefixCoder[FastPrimitivesCoder],
IterableCoder[LengthPrefixCoder[FastPrimitivesCoder]]]], len(consumers)=1]]>
DEBUG:apache_beam.io.gcp.bigquery:Creating or getting table <TableReference
datasetId: 'python_query_to_table_15843762783322'
projectId: 'apache-beam-testing'
tableId: 'output_table'> with schema {'fields': [{'type': u'STRING', 'name':
u'fruit', 'mode': 'NULLABLE'}]}.
DEBUG:apache_beam.io.gcp.bigquery_tools:Created the table with id output_table
INFO:apache_beam.io.gcp.bigquery_tools:Created table
apache-beam-testing.python_query_to_table_15843762783322.output_table with
schema <TableSchema
fields: [<TableFieldSchema
fields: []
mode: u'NULLABLE'
name: u'fruit'
type: u'STRING'>]>. Result: <Table
creationTime: 1584376293486
etag: u'D2eOZIfnoQqQ8e9GRKlilw=='
id: u'apache-beam-testing:python_query_to_table_15843762783322.output_table'
kind: u'bigquery#table'
lastModifiedTime: 1584376293531
location: u'US'
numBytes: 0
numLongTermBytes: 0
numRows: 0
schema: <TableSchema
fields: [<TableFieldSchema
fields: []
mode: u'NULLABLE'
name: u'fruit'
type: u'STRING'>]>
selfLink:
u'https://bigquery.googleapis.com/bigquery/v2/projects/apache-beam-testing/datasets/python_query_to_table_15843762783322/tables/output_table'
tableReference: <TableReference
datasetId: u'python_query_to_table_15843762783322'
projectId: u'apache-beam-testing'
tableId: u'output_table'>
type: u'TABLE'>.
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DataInputOperation
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read
receivers=[SingletonConsumerSet[write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/GroupByKey/Read.out0,
coder=WindowedValueCoder[TupleCoder[LengthPrefixCoder[FastPrimitivesCoder],
IterableCoder[LengthPrefixCoder[FastPrimitivesCoder]]]], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DoOperation
write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps)
output_tags=[u'None'],
receivers=[SingletonConsumerSet[write/_StreamToBigQuery/CommitInsertIds/ReshufflePerKey/FlatMap(restore_timestamps).out0,
coder=WindowedValueCoder[TupleCoder[FastPrimitivesCoder,
TupleCoder[FastPrimitivesCoder, TupleCoder[FastPrimitivesCoder,
FastPrimitivesCoder]]]], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DoOperation
write/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys output_tags=[u'None'],
receivers=[SingletonConsumerSet[write/_StreamToBigQuery/CommitInsertIds/RemoveRandomKeys.out0,
coder=WindowedValueCoder[TupleCoder[FastPrimitivesCoder,
TupleCoder[FastPrimitivesCoder, FastPrimitivesCoder]]], len(consumers)=1]]>
DEBUG:apache_beam.runners.worker.bundle_processor:finish <DoOperation
write/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn)
output_tags=[u'None', u'FailedRows'],
receivers=[ConsumerSet[write/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn).out0,
coder=WindowedValueCoder[FastPrimitivesCoder], len(consumers)=0],
ConsumerSet[write/_StreamToBigQuery/StreamInsertRows/ParDo(BigQueryWriteFn).out1,
coder=WindowedValueCoder[FastPrimitivesCoder], len(consumers)=0]]>
DEBUG:apache_beam.io.gcp.bigquery:Attempting to flush to all destinations.
Total buffered: 2
DEBUG:apache_beam.io.gcp.bigquery:Flushing data to
apache-beam-testing:python_query_to_table_15843762783322.output_table. Total 2
rows.
DEBUG:apache_beam.runners.portability.fn_api_runner:Wait for the bundle
bundle_12 to finish.
INFO:apache_beam.io.gcp.tests.bigquery_matcher:Attempting to perform query
SELECT fruit from `python_query_to_table_15843762783322.output_table`; to BQ
DEBUG:google.auth.transport._http_client:Making request: GET
http://169.254.169.254
DEBUG:google.auth.transport._http_client:Making request: GET
http://metadata.google.internal/computeMetadata/v1/project/project-id
DEBUG:urllib3.util.retry:Converted retries value: 3 -> Retry(total=3,
connect=None, read=None, redirect=None, status=None)
DEBUG:google.auth.transport.requests:Making request: GET
http://metadata.google.internal/computeMetadata/v1/instance/service-accounts/default/?recursive=true
DEBUG:urllib3.connectionpool:Starting new HTTP connection (1):
metadata.google.internal:80
DEBUG:urllib3.connectionpool:http://metadata.google.internal:80 "GET
/computeMetadata/v1/instance/service-accounts/default/?recursive=true HTTP/1.1"
200 144
DEBUG:google.auth.transport.requests:Making request: GET
http://metadata.google.internal/computeMetadata/v1/instance/service-accounts/[email protected]/token
DEBUG:urllib3.connectionpool:http://metadata.google.internal:80 "GET
/computeMetadata/v1/instance/service-accounts/[email protected]/token
HTTP/1.1" 200 192
DEBUG:urllib3.connectionpool:Starting new HTTPS connection (1):
bigquery.googleapis.com:443
DEBUG:urllib3.connectionpool:https://bigquery.googleapis.com:443 "POST
/bigquery/v2/projects/apache-beam-testing/jobs HTTP/1.1" 200 None
DEBUG:urllib3.connectionpool:https://bigquery.googleapis.com:443 "GET
/bigquery/v2/projects/apache-beam-testing/queries/1bf5ed04-f0ca-4bac-b10c-07d5ac63f304?location=US&maxResults=0
HTTP/1.1" 200 None
DEBUG:urllib3.connectionpool:https://bigquery.googleapis.com:443 "GET
/bigquery/v2/projects/apache-beam-testing/datasets/_7357fab0f784d2a7327ddbe81cdd1f4ca7e429cd/tables/anone87e0c70_3401_456b_aecc_161191f61d99/data
HTTP/1.1" 200 None
INFO:apache_beam.io.gcp.tests.bigquery_matcher:Read from given query (SELECT
fruit from `python_query_to_table_15843762783322.output_table`;), total rows 2
INFO:apache_beam.io.gcp.tests.bigquery_matcher:Generate checksum:
3b2cefe89863bf492d48f7d4da960f2999802a89
test_big_query_legacy_sql
(apache_beam.io.gcp.big_query_query_to_table_it_test.BigQueryQueryToTableIT)
... ok
test_big_query_new_types
(apache_beam.io.gcp.big_query_query_to_table_it_test.BigQueryQueryToTableIT)
... ok
test_big_query_new_types_native
(apache_beam.io.gcp.big_query_query_to_table_it_test.BigQueryQueryToTableIT)
... ok
test_big_query_standard_sql
(apache_beam.io.gcp.big_query_query_to_table_it_test.BigQueryQueryToTableIT)
... ok
test_big_query_standard_sql_kms_key_native
(apache_beam.io.gcp.big_query_query_to_table_it_test.BigQueryQueryToTableIT)
... SKIP: This test doesn't work on DirectRunner.
----------------------------------------------------------------------
XML: nosetests-directRunnerIT-batch.xml
----------------------------------------------------------------------
XML:
<https://builds.apache.org/job/beam_PostCommit_Python2/ws/src/sdks/python/nosetests.xml>
----------------------------------------------------------------------
Ran 19 tests in 57.176s
OK (SKIP=1)
>>> RUNNING integration tests with pipeline options: --runner=TestDirectRunner
>>> --project=apache-beam-testing
>>> --staging_location=gs://temp-storage-for-end-to-end-tests/staging-it
>>> --temp_location=gs://temp-storage-for-end-to-end-tests/temp-it
>>> --output=gs://temp-storage-for-end-to-end-tests/py-it-cloud/output
>>> --sdk_location=build/apache-beam.tar.gz
>>> --requirements_file=postcommit_requirements.txt --num_workers=1
>>> --sleep_secs=20 --streaming
>>> --kms_key_name=projects/apache-beam-testing/locations/global/keyRings/beam-it/cryptoKeys/test
>>>
>>> --dataflow_kms_key=projects/apache-beam-testing/locations/global/keyRings/beam-it/cryptoKeys/test
>>> test options: --nocapture --processes=8 --process-timeout=4500
>>> --tests=apache_beam.examples.wordcount_it_test:WordCountIT.test_wordcount_it,apache_beam.io.gcp.pubsub_integration_test:PubSubIntegrationTest,apache_beam.io.gcp.bigquery_test:BigQueryStreamingInsertTransformIntegrationTests.test_multiple_destinations_transform,apache_beam.io.gcp.bigquery_test:PubSubBigQueryIT,apache_beam.io.gcp.bigquery_file_loads_test:BigQueryFileLoadsIT.test_bqfl_streaming
setup.py:251: UserWarning: You are using Apache Beam with Python 2. New
releases of Apache Beam will soon support Python 3 only.
'You are using Apache Beam with Python 2. '
<https://builds.apache.org/job/beam_PostCommit_Python2/ws/src/build/gradleenv/45127155/local/lib/python2.7/site-packages/setuptools/dist.py>:476:
UserWarning: Normalizing '2.21.0.dev' to '2.21.0.dev0'
normalized_version,
running nosetests
running egg_info
INFO:gen_protos:Skipping proto regeneration: all files up to date
writing requirements to apache_beam.egg-info/requires.txt
writing apache_beam.egg-info/PKG-INFO
writing top-level names to apache_beam.egg-info/top_level.txt
writing dependency_links to apache_beam.egg-info/dependency_links.txt
writing entry points to apache_beam.egg-info/entry_points.txt
reading manifest file 'apache_beam.egg-info/SOURCES.txt'
reading manifest template 'MANIFEST.in'
warning: no files found matching 'README.md'
warning: no files found matching 'NOTICE'
warning: no files found matching 'LICENSE'
writing manifest file 'apache_beam.egg-info/SOURCES.txt'
<https://builds.apache.org/job/beam_PostCommit_Python2/ws/src/sdks/python/apache_beam/__init__.py>:82:
UserWarning: You are using Apache Beam with Python 2. New releases of Apache
Beam will soon support Python 3 only.
'You are using Apache Beam with Python 2. '
INFO:root:Generating grammar tables from /usr/lib/python2.7/lib2to3/Grammar.txt
INFO:root:Generating grammar tables from
/usr/lib/python2.7/lib2to3/PatternGrammar.txt
> Task :sdks:python:test-suites:dataflow:py2:postCommitIT
Traceback (most recent call last):
File "setup.py", line 315, in <module>
'mypy': generate_protos_first(mypy),
File
"<https://builds.apache.org/job/beam_PostCommit_Python2/ws/src/build/gradleenv/-194514014/local/lib/python2.7/site-packages/setuptools/__init__.py",>
line 145, in setup
return distutils.core.setup(**attrs)
File "/usr/lib/python2.7/distutils/core.py", line 151, in setup
dist.run_commands()
File "/usr/lib/python2.7/distutils/dist.py", line 953, in run_commands
self.run_command(cmd)
File "/usr/lib/python2.7/distutils/dist.py", line 972, in run_command
cmd_obj.run()
File
"<https://builds.apache.org/job/beam_PostCommit_Python2/ws/src/build/gradleenv/-194514014/local/lib/python2.7/site-packages/nose/commands.py",>
line 158, in run
TestProgram(argv=argv, config=self.__config)
File
"<https://builds.apache.org/job/beam_PostCommit_Python2/ws/src/build/gradleenv/-194514014/local/lib/python2.7/site-packages/nose/core.py",>
line 121, in __init__
**extra_args)
File "/usr/lib/python2.7/unittest/main.py", line 95, in __init__
self.runTests()
File
"<https://builds.apache.org/job/beam_PostCommit_Python2/ws/src/build/gradleenv/-194514014/local/lib/python2.7/site-packages/nose/core.py",>
line 207, in runTests
result = self.testRunner.run(self.test)
File
"<https://builds.apache.org/job/beam_PostCommit_Python2/ws/src/build/gradleenv/-194514014/local/lib/python2.7/site-packages/nose/plugins/multiprocess.py",>
line 396, in run
timeout=nexttimeout)
File "<string>", line 2, in get
File "/usr/lib/python2.7/multiprocessing/managers.py", line 759, in
_callmethod
Terminated
> Task :sdks:python:test-suites:dataflow:py2:postCommitIT FAILED
Daemon will be stopped at the end of the build after the daemon was no longer
found in the daemon registry
The message received from the daemon indicates that the daemon has disappeared.
Build request sent: Build{id=9a05e247-b1d9-4edd-8b4b-d6834026e739,
currentDir=<https://builds.apache.org/job/beam_PostCommit_Python2/ws/src}>
Attempting to read last messages from the daemon log...
Daemon pid: 13021
log file: /home/jenkins/.gradle/daemon/5.2.1/daemon-13021.out.log
----- Last 20 lines from daemon log file - daemon-13021.out.log -----
at
org.gradle.process.internal.DefaultExecHandle.execExceptionFor(DefaultExecHandle.java:232)
at
org.gradle.process.internal.DefaultExecHandle.setEndStateInfo(DefaultExecHandle.java:209)
at
org.gradle.process.internal.DefaultExecHandle.failed(DefaultExecHandle.java:356)
at
org.gradle.process.internal.ExecHandleRunner.run(ExecHandleRunner.java:86)
at
org.gradle.internal.operations.CurrentBuildOperationPreservingRunnable.run(CurrentBuildOperationPreservingRunnable.java:42)
at
org.gradle.internal.concurrent.ExecutorPolicy$CatchAndRecordFailures.onExecute(ExecutorPolicy.java:63)
at
org.gradle.internal.concurrent.ManagedExecutorImpl$1.run(ManagedExecutorImpl.java:46)
at
java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at
java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at
org.gradle.internal.concurrent.ThreadFactoryImpl$ManagedThreadRunnable.run(ThreadFactoryImpl.java:55)
at java.lang.Thread.run(Thread.java:748)
Caused by: java.lang.IllegalStateException: Shutdown in progress
at
java.lang.ApplicationShutdownHooks.remove(ApplicationShutdownHooks.java:82)
at java.lang.Runtime.removeShutdownHook(Runtime.java:240)
at
org.gradle.process.internal.shutdown.ShutdownHooks.removeShutdownHook(ShutdownHooks.java:33)
at
org.gradle.process.internal.DefaultExecHandle.setEndStateInfo(DefaultExecHandle.java:199)
at
org.gradle.process.internal.DefaultExecHandle.aborted(DefaultExecHandle.java:352)
at
org.gradle.process.internal.ExecHandleRunner.completed(ExecHandleRunner.java:107)
at
org.gradle.process.internal.ExecHandleRunner.run(ExecHandleRunner.java:83)
... 7 more
----- End of the daemon log -----
FAILURE: Build failed with an exception.
* What went wrong:
Gradle build daemon disappeared unexpectedly (it may have been killed or may
have crashed)
* Try:
Run with --stacktrace option to get the stack trace. Run with --info or --debug
option to get more log output. Run with --scan to get full insights.
* Get more help at https://help.gradle.org
Build step 'Invoke Gradle script' changed build result to FAILURE
Build step 'Invoke Gradle script' marked build as failure
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]