Task:
('hash-join-4b4b25fcaf703c1e99e652bac3539b7f', 3)
Exception:
RuntimeError('Worker tcp://172.21.25.29:34733 left during active shuffle 3d4472ea35c0f53629acee4ad7bde540')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 902, in get_output_partition
shuffle = self.get_shuffle_run(shuffle_id, run_id)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 873, in get_shuffle_run
return sync(
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 678, in _get_shuffle_run
shuffle = await self._refresh_shuffle(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 780, in _refresh_shuffle
raise RuntimeError(result["message"])
Task:
('shuffle-p2p-210aea3b503452ad401d37bd774444eb', 11)
Exception:
RuntimeError('shuffle_unpack failed during shuffle 416a4a8a9303bd1b6a928557fd2c4409')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 86, in shuffle_unpack
raise RuntimeError(msg)
Task:
('shuffle-p2p-210aea3b503452ad401d37bd774444eb', 0)
Exception:
RuntimeError('shuffle_unpack failed during shuffle 416a4a8a9303bd1b6a928557fd2c4409')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 86, in shuffle_unpack
raise RuntimeError(msg)
Task:
('hash-join-transfer-2f4fca06bc9c5475b26c12cf93722f4f', 0)
Exception:
RuntimeError('shuffle_transfer failed during shuffle 2f4fca06bc9c5475b26c12cf93722f4f')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 145, in merge_transfer
return shuffle_transfer(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 71, in shuffle_transfer
raise RuntimeError(msg)
Task:
run-13177c827a5279df9faaca249711e262
Exception:
KilledWorker("('from_pandas-7257ad7e7e1317b17a9f46f3bd9f975f', 10)", <WorkerState 'tcp://172.21.159.246:38983', status: closed, memory: 0, processing: 0>, 3)
File "/opt/conda/lib/python3.10/site-packages/algorithms/anomalies_detection/anomalies_increase_and_decrease.py", line 295, in run
if len(cam_data_dask) == 0:
File "/opt/conda/lib/python3.10/site-packages/dask/dataframe/core.py", line 4775, in __len__
return len(s)
File "/opt/conda/lib/python3.10/site-packages/dask/dataframe/core.py", line 843, in __len__
).compute()
File "/opt/conda/lib/python3.10/site-packages/dask/base.py", line 314, in compute
(result,) = compute(self, traverse=False, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/dask/base.py", line 599, in compute
results = schedule(dsk, keys, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/distributed/client.py", line 3186, in get
results = self.gather(packed, asynchronous=asynchronous, direct=direct)
File "/opt/conda/lib/python3.10/site-packages/distributed/client.py", line 2345, in gather
return self.sync(
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 349, in sync
return sync(
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/client.py", line 2208, in _gather
raise exception.with_traceback(traceback)
Task:
('assign-f675a204035f907c92325171b70e20d1', 1)
Exception:
GEOSException('IllegalArgumentException: point array must contain 0 or >1 elements\n')
File "/opt/conda/lib/python3.10/site-packages/dask/optimization.py", line 990, in __call__
return core.get(self.dsk, self.outkey, dict(zip(self.inkeys, args)))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 149, in get
result = _execute_task(task, cache)
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 113, in _execute_task
return [_execute_task(a, cache) for a in arg]
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 113, in <listcomp>
return [_execute_task(a, cache) for a in arg]
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/utils.py", line 73, in apply
return func(*args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/dask/dataframe/core.py", line 7006, in apply_and_enforce
df = func(*args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/dask/dataframe/groupby.py", line 230, in _groupby_slice_apply
return g.apply(func, *args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/pandas/core/groupby/generic.py", line 254, in apply
return super().apply(func, *args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/pandas/core/groupby/groupby.py", line 1567, in apply
result = self._python_apply_general(f, self._selected_obj)
File "/opt/conda/lib/python3.10/site-packages/pandas/core/groupby/groupby.py", line 1629, in _python_apply_general
values, mutated = self.grouper.apply(f, data, self.axis)
File "/opt/conda/lib/python3.10/site-packages/pandas/core/groupby/ops.py", line 839, in apply
res = f(group)
File "/app/vehicle_distance_metric_worker.py", line 43, in <lambda>
File "/opt/conda/lib/python3.10/site-packages/shapely/geometry/linestring.py", line 73, in __new__
geom = shapely.linestrings(coordinates)
File "/opt/conda/lib/python3.10/site-packages/shapely/decorators.py", line 77, in wrapped
return func(*args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/shapely/creation.py", line 120, in linestrings
return lib.linestrings(coords, out=out, **kwargs)
Task:
('hash-join-transfer-05b2eea45734abca7b550c5b98f8f3a6', 13)
Exception:
RuntimeError('shuffle_transfer failed during shuffle 05b2eea45734abca7b550c5b98f8f3a6')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 145, in merge_transfer
return shuffle_transfer(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 71, in shuffle_transfer
raise RuntimeError(msg)
Task:
('hash-join-transfer-eec3c2cf83943d5b1a0c93ed88834f19', 9)
Exception:
RuntimeError('shuffle_transfer failed during shuffle eec3c2cf83943d5b1a0c93ed88834f19')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 145, in merge_transfer
return shuffle_transfer(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 71, in shuffle_transfer
raise RuntimeError(msg)
Task:
('hash-join-3c01ff980d16020c70c3d0faede248f4', 9)
Exception:
ArrowInvalid('Schema at index 1 was different: \nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: string\nlongitude: double\nlatitude: double\naltitude: double\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]\nvs\nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: null\nlongitude: double\nlatitude: double\naltitude: double\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 903, in get_output_partition
return sync(self.worker.loop, shuffle.get_output_partition, output_partition)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 529, in get_output_partition
out = await self.offload(_)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 132, in offload
return await asyncio.get_running_loop().run_in_executor(
File "/opt/conda/lib/python3.10/concurrent/futures/thread.py", line 58, in run
result = self.fn(*self.args, **self.kwargs)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 526, in _
df = convert_partition(data)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_arrow.py", line 57, in convert_partition
return pa.concat_tables(shards)
File "pyarrow/table.pxi", line 5224, in pyarrow.lib.concat_tables
File "pyarrow/error.pxi", line 144, in pyarrow.lib.pyarrow_internal_check_status
File "pyarrow/error.pxi", line 100, in pyarrow.lib.check_status
Task:
('shuffle-p2p-5444bd365d67f74c247fc292b1eaa1f9', 9)
Exception:
RuntimeError('shuffle_unpack failed during shuffle 8f860c6fe1f949e8bb18fe310bb37739')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 86, in shuffle_unpack
raise RuntimeError(msg)
Task:
('shuffle-p2p-5444bd365d67f74c247fc292b1eaa1f9', 15)
Exception:
RuntimeError('shuffle_unpack failed during shuffle 8f860c6fe1f949e8bb18fe310bb37739')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 86, in shuffle_unpack
raise RuntimeError(msg)
Task:
('shuffle-p2p-084701763c90f362841bf1fc742c3146', 4)
Exception:
RuntimeError('shuffle_unpack failed during shuffle 080fc1af058f130722ca3ee502f17ab5')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 86, in shuffle_unpack
raise RuntimeError(msg)
Task:
('hash-join-703532e48cd395241db47262dd0c1962', 2)
Exception:
ArrowInvalid('Schema at index 1 was different: \nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: null\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]\nvs\nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: string\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 903, in get_output_partition
return sync(self.worker.loop, shuffle.get_output_partition, output_partition)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 529, in get_output_partition
out = await self.offload(_)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 132, in offload
return await asyncio.get_running_loop().run_in_executor(
File "/opt/conda/lib/python3.10/concurrent/futures/thread.py", line 58, in run
result = self.fn(*self.args, **self.kwargs)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 526, in _
df = convert_partition(data)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_arrow.py", line 57, in convert_partition
return pa.concat_tables(shards)
File "pyarrow/table.pxi", line 5224, in pyarrow.lib.concat_tables
File "pyarrow/error.pxi", line 144, in pyarrow.lib.pyarrow_internal_check_status
File "pyarrow/error.pxi", line 100, in pyarrow.lib.check_status
Task:
('lambda-f692a613a077480174270702ec5eff83', 13)
Exception:
ValueError('Found array with 1 sample(s) (shape=(1, 1)) while a minimum of 2 is required by AgglomerativeClustering.')
File "/opt/conda/lib/python3.10/site-packages/dask/optimization.py", line 990, in __call__
return core.get(self.dsk, self.outkey, dict(zip(self.inkeys, args)))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 149, in get
result = _execute_task(task, cache)
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/utils.py", line 73, in apply
return func(*args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/dask/dataframe/core.py", line 7006, in apply_and_enforce
df = func(*args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/dask/dataframe/groupby.py", line 230, in _groupby_slice_apply
return g.apply(func, *args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/pandas/core/groupby/groupby.py", line 1567, in apply
result = self._python_apply_general(f, self._selected_obj)
File "/opt/conda/lib/python3.10/site-packages/pandas/core/groupby/groupby.py", line 1629, in _python_apply_general
values, mutated = self.grouper.apply(f, data, self.axis)
File "/opt/conda/lib/python3.10/site-packages/pandas/core/groupby/ops.py", line 839, in apply
res = f(group)
File "/app/special_vehicle_stop_metric_worker.py", line 48, in <lambda>
File "/opt/conda/lib/python3.10/site-packages/sklearn/cluster/_agglomerative.py", line 1099, in fit_predict
return super().fit_predict(X, y)
File "/opt/conda/lib/python3.10/site-packages/sklearn/base.py", line 753, in fit_predict
self.fit(X)
File "/opt/conda/lib/python3.10/site-packages/sklearn/cluster/_agglomerative.py", line 955, in fit
X = self._validate_data(X, ensure_min_samples=2)
File "/opt/conda/lib/python3.10/site-packages/sklearn/base.py", line 565, in _validate_data
X = check_array(X, input_name="X", **check_params)
File "/opt/conda/lib/python3.10/site-packages/sklearn/utils/validation.py", line 931, in check_array
raise ValueError(
Task:
('hash-join-0ad9eedf603cf4e330cfc753a28d9c23', 9)
Exception:
ArrowInvalid('Schema at index 1 was different: \nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: null\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]\nvs\nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: string\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 903, in get_output_partition
return sync(self.worker.loop, shuffle.get_output_partition, output_partition)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 529, in get_output_partition
out = await self.offload(_)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 132, in offload
return await asyncio.get_running_loop().run_in_executor(
File "/opt/conda/lib/python3.10/concurrent/futures/thread.py", line 58, in run
result = self.fn(*self.args, **self.kwargs)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 526, in _
df = convert_partition(data)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_arrow.py", line 57, in convert_partition
return pa.concat_tables(shards)
File "pyarrow/table.pxi", line 5224, in pyarrow.lib.concat_tables
File "pyarrow/error.pxi", line 144, in pyarrow.lib.pyarrow_internal_check_status
File "pyarrow/error.pxi", line 100, in pyarrow.lib.check_status
Task:
('hash-join-fdf4a6ab751c933e84ff5c171d672dc8', 14)
Exception:
RuntimeError('Worker tcp://172.21.151.150:38387 left during active shuffle 7fe06a3da8052015aabaf82be4107670')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 902, in get_output_partition
shuffle = self.get_shuffle_run(shuffle_id, run_id)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 873, in get_shuffle_run
return sync(
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 678, in _get_shuffle_run
shuffle = await self._refresh_shuffle(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 780, in _refresh_shuffle
raise RuntimeError(result["message"])
Task:
('shuffle-p2p-ee7354768e111c1f5f98db89d91d0ab5', 7)
Exception:
RuntimeError('shuffle_unpack failed during shuffle 14b747bcff3a62b3bed39aaa4dfdaa24')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 86, in shuffle_unpack
raise RuntimeError(msg)
Task:
('shuffle-p2p-ee7354768e111c1f5f98db89d91d0ab5', 6)
Exception:
RuntimeError('shuffle_unpack failed during shuffle 14b747bcff3a62b3bed39aaa4dfdaa24')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 86, in shuffle_unpack
raise RuntimeError(msg)
Task:
('hash-join-181407cb727510482ffc7872df4af145', 6)
Exception:
ArrowInvalid('Schema at index 1 was different: \nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: string\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]\nvs\nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: null\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 903, in get_output_partition
return sync(self.worker.loop, shuffle.get_output_partition, output_partition)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 529, in get_output_partition
out = await self.offload(_)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 132, in offload
return await asyncio.get_running_loop().run_in_executor(
File "/opt/conda/lib/python3.10/concurrent/futures/thread.py", line 58, in run
result = self.fn(*self.args, **self.kwargs)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 526, in _
df = convert_partition(data)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_arrow.py", line 57, in convert_partition
return pa.concat_tables(shards)
File "pyarrow/table.pxi", line 5224, in pyarrow.lib.concat_tables
File "pyarrow/error.pxi", line 144, in pyarrow.lib.pyarrow_internal_check_status
File "pyarrow/error.pxi", line 100, in pyarrow.lib.check_status
Task:
('assign-fbe758a3a0fb0195da07b7b12385cdbf', 1)
Exception:
GEOSException('IllegalArgumentException: point array must contain 0 or >1 elements\n')
File "/opt/conda/lib/python3.10/site-packages/dask/optimization.py", line 990, in __call__
return core.get(self.dsk, self.outkey, dict(zip(self.inkeys, args)))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 149, in get
result = _execute_task(task, cache)
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 113, in _execute_task
return [_execute_task(a, cache) for a in arg]
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 113, in <listcomp>
return [_execute_task(a, cache) for a in arg]
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/utils.py", line 73, in apply
return func(*args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/dask/dataframe/core.py", line 7006, in apply_and_enforce
df = func(*args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/dask/dataframe/groupby.py", line 230, in _groupby_slice_apply
return g.apply(func, *args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/pandas/core/groupby/generic.py", line 254, in apply
return super().apply(func, *args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/pandas/core/groupby/groupby.py", line 1567, in apply
result = self._python_apply_general(f, self._selected_obj)
File "/opt/conda/lib/python3.10/site-packages/pandas/core/groupby/groupby.py", line 1629, in _python_apply_general
values, mutated = self.grouper.apply(f, data, self.axis)
File "/opt/conda/lib/python3.10/site-packages/pandas/core/groupby/ops.py", line 839, in apply
res = f(group)
File "/app/vehicle_distance_metric_worker.py", line 43, in <lambda>
File "/opt/conda/lib/python3.10/site-packages/shapely/geometry/linestring.py", line 73, in __new__
geom = shapely.linestrings(coordinates)
File "/opt/conda/lib/python3.10/site-packages/shapely/decorators.py", line 77, in wrapped
return func(*args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/shapely/creation.py", line 120, in linestrings
return lib.linestrings(coords, out=out, **kwargs)
Task:
('hash-join-a80329184eccb02b1f490313add46065', 5)
Exception:
ArrowInvalid('Schema at index 1 was different: \nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: string\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]\nvs\nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: null\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 903, in get_output_partition
return sync(self.worker.loop, shuffle.get_output_partition, output_partition)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 529, in get_output_partition
out = await self.offload(_)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 132, in offload
return await asyncio.get_running_loop().run_in_executor(
File "/opt/conda/lib/python3.10/concurrent/futures/thread.py", line 58, in run
result = self.fn(*self.args, **self.kwargs)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 526, in _
df = convert_partition(data)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_arrow.py", line 57, in convert_partition
return pa.concat_tables(shards)
File "pyarrow/table.pxi", line 5224, in pyarrow.lib.concat_tables
File "pyarrow/error.pxi", line 144, in pyarrow.lib.pyarrow_internal_check_status
File "pyarrow/error.pxi", line 100, in pyarrow.lib.check_status
Task:
('hash-join-24a31f1dadf8bf09023eb1d3e9e24c90', 0)
Exception:
RuntimeError('Worker tcp://172.21.151.150:44549 left during active shuffle 40a265a636f672d90a667506167eb84d')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 902, in get_output_partition
shuffle = self.get_shuffle_run(shuffle_id, run_id)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 873, in get_shuffle_run
return sync(
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 678, in _get_shuffle_run
shuffle = await self._refresh_shuffle(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 780, in _refresh_shuffle
raise RuntimeError(result["message"])
Task:
('hash-join-042976186e2fb73429ddab1c1f4a0c61', 2)
Exception:
RuntimeError('Worker tcp://172.21.159.246:36409 left during active shuffle 8135f54f9dce30185f4d3fef2ada6bdd')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 902, in get_output_partition
shuffle = self.get_shuffle_run(shuffle_id, run_id)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 873, in get_shuffle_run
return sync(
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 678, in _get_shuffle_run
shuffle = await self._refresh_shuffle(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 780, in _refresh_shuffle
raise RuntimeError(result["message"])
Task:
('shuffle-transfer-2c7fd55cad0e13fedbad3273e327127d', 14)
Exception:
RuntimeError('shuffle_transfer failed during shuffle 2c7fd55cad0e13fedbad3273e327127d')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 71, in shuffle_transfer
raise RuntimeError(msg)
Task:
shuffle-barrier-4525ea1da2d0e9579776be9fe4301954
Exception:
RuntimeError('shuffle_barrier failed during shuffle 4525ea1da2d0e9579776be9fe4301954')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 97, in shuffle_barrier
raise RuntimeError(msg)
Task:
('assign-3361cddacb9ac6054e95a968f5d5153a', 6)
Exception:
GEOSException('IllegalArgumentException: point array must contain 0 or >1 elements\n')
File "/opt/conda/lib/python3.10/site-packages/dask/optimization.py", line 990, in __call__
return core.get(self.dsk, self.outkey, dict(zip(self.inkeys, args)))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 149, in get
result = _execute_task(task, cache)
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr>
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 113, in _execute_task
return [_execute_task(a, cache) for a in arg]
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 113, in <listcomp>
return [_execute_task(a, cache) for a in arg]
File "/opt/conda/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task
return func(*(_execute_task(a, cache) for a in args))
File "/opt/conda/lib/python3.10/site-packages/dask/utils.py", line 73, in apply
return func(*args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/dask/dataframe/core.py", line 7006, in apply_and_enforce
df = func(*args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/dask/dataframe/groupby.py", line 230, in _groupby_slice_apply
return g.apply(func, *args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/pandas/core/groupby/generic.py", line 254, in apply
return super().apply(func, *args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/pandas/core/groupby/groupby.py", line 1567, in apply
result = self._python_apply_general(f, self._selected_obj)
File "/opt/conda/lib/python3.10/site-packages/pandas/core/groupby/groupby.py", line 1629, in _python_apply_general
values, mutated = self.grouper.apply(f, data, self.axis)
File "/opt/conda/lib/python3.10/site-packages/pandas/core/groupby/ops.py", line 839, in apply
res = f(group)
File "/app/vehicle_distance_metric_worker.py", line 43, in <lambda>
File "/opt/conda/lib/python3.10/site-packages/shapely/geometry/linestring.py", line 73, in __new__
geom = shapely.linestrings(coordinates)
File "/opt/conda/lib/python3.10/site-packages/shapely/decorators.py", line 77, in wrapped
return func(*args, **kwargs)
File "/opt/conda/lib/python3.10/site-packages/shapely/creation.py", line 120, in linestrings
return lib.linestrings(coords, out=out, **kwargs)
Task:
('hash-join-c89798abd68c8b783c220d4c4b8fe23b', 15)
Exception:
RuntimeError('Worker tcp://172.21.25.54:39899 left during active shuffle 64ad8ad5105af031ddd958c4a37785f7')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 902, in get_output_partition
shuffle = self.get_shuffle_run(shuffle_id, run_id)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 873, in get_shuffle_run
return sync(
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 678, in _get_shuffle_run
shuffle = await self._refresh_shuffle(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 780, in _refresh_shuffle
raise RuntimeError(result["message"])
Task:
('shuffle-p2p-dfa8a71f87f499add737ab6f7384efc4', 1)
Exception:
RuntimeError('shuffle_unpack failed during shuffle e3ef17c11773502e54c27ab45b7c658a')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 86, in shuffle_unpack
raise RuntimeError(msg)
Task:
('hash-join-569f92d5b293d53137f7a28d5370cbbf', 8)
Exception:
ArrowInvalid('Schema at index 1 was different: \nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: null\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]\nvs\nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: string\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 903, in get_output_partition
return sync(self.worker.loop, shuffle.get_output_partition, output_partition)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 529, in get_output_partition
out = await self.offload(_)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 132, in offload
return await asyncio.get_running_loop().run_in_executor(
File "/opt/conda/lib/python3.10/concurrent/futures/thread.py", line 58, in run
result = self.fn(*self.args, **self.kwargs)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 526, in _
df = convert_partition(data)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_arrow.py", line 57, in convert_partition
return pa.concat_tables(shards)
File "pyarrow/table.pxi", line 5224, in pyarrow.lib.concat_tables
File "pyarrow/error.pxi", line 144, in pyarrow.lib.pyarrow_internal_check_status
File "pyarrow/error.pxi", line 100, in pyarrow.lib.check_status
Task:
('shuffle-p2p-3e35f09f42cc44da10ab45ef7aa3ab9e', 15)
Exception:
RuntimeError('shuffle_unpack failed during shuffle ce71486c0e31fb151f9ca45582783a7f')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 86, in shuffle_unpack
raise RuntimeError(msg)
Task:
('hash-join-5f46ad7d5092a6fdc1641a7f17d692ee', 4)
Exception:
ArrowInvalid('Schema at index 1 was different: \nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: string\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]\nvs\nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: null\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 903, in get_output_partition
return sync(self.worker.loop, shuffle.get_output_partition, output_partition)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 529, in get_output_partition
out = await self.offload(_)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 132, in offload
return await asyncio.get_running_loop().run_in_executor(
File "/opt/conda/lib/python3.10/concurrent/futures/thread.py", line 58, in run
result = self.fn(*self.args, **self.kwargs)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 526, in _
df = convert_partition(data)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_arrow.py", line 57, in convert_partition
return pa.concat_tables(shards)
File "pyarrow/table.pxi", line 5224, in pyarrow.lib.concat_tables
File "pyarrow/error.pxi", line 144, in pyarrow.lib.pyarrow_internal_check_status
File "pyarrow/error.pxi", line 100, in pyarrow.lib.check_status
Task:
('hash-join-db0706fcfd504ffb297ed7c0b4f5c5b1', 10)
Exception:
RuntimeError('Worker tcp://172.21.159.253:40571 left during active shuffle 9a400d2158ab3d046210d22e656072bc')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 902, in get_output_partition
shuffle = self.get_shuffle_run(shuffle_id, run_id)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 873, in get_shuffle_run
return sync(
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 678, in _get_shuffle_run
shuffle = await self._refresh_shuffle(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 780, in _refresh_shuffle
raise RuntimeError(result["message"])
Task:
('hash-join-transfer-ee6fa7f2d2572fa5840c05b4e74410bd', 5)
Exception:
RuntimeError('shuffle_transfer failed during shuffle ee6fa7f2d2572fa5840c05b4e74410bd')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 145, in merge_transfer
return shuffle_transfer(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 71, in shuffle_transfer
raise RuntimeError(msg)
Task:
shuffle-barrier-011f1efefe4c945109110a4b577cbf4d
Exception:
RuntimeError('shuffle_barrier failed during shuffle 011f1efefe4c945109110a4b577cbf4d')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 97, in shuffle_barrier
raise RuntimeError(msg)
Task:
('shuffle-transfer-b09e00533ac31cc9c2c310d2c2305e75', 15)
Exception:
RuntimeError('shuffle_transfer failed during shuffle b09e00533ac31cc9c2c310d2c2305e75')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 71, in shuffle_transfer
raise RuntimeError(msg)
Task:
('hash-join-2d542f68295535bbdcb65796faa99453', 10)
Exception:
RuntimeError('Worker tcp://172.21.159.222:35865 left during active shuffle 97cc0e65355bd18db884b23ece825ae9')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 902, in get_output_partition
shuffle = self.get_shuffle_run(shuffle_id, run_id)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 873, in get_shuffle_run
return sync(
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 678, in _get_shuffle_run
shuffle = await self._refresh_shuffle(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 780, in _refresh_shuffle
raise RuntimeError(result["message"])
Task:
('hash-join-02e2b224aac843df023a3c0c3f6656d4', 13)
Exception:
RuntimeError('Worker tcp://172.21.159.207:43461 left during active shuffle e90d5c7caa33f2f3560741d3a1f64fab')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 902, in get_output_partition
shuffle = self.get_shuffle_run(shuffle_id, run_id)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 873, in get_shuffle_run
return sync(
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 678, in _get_shuffle_run
shuffle = await self._refresh_shuffle(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 780, in _refresh_shuffle
raise RuntimeError(result["message"])
Task:
shuffle-barrier-a079868710da4a8af2c4f117f9c314dc
Exception:
RuntimeError('shuffle_barrier failed during shuffle a079868710da4a8af2c4f117f9c314dc')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 97, in shuffle_barrier
raise RuntimeError(msg)
Task:
('shuffle-p2p-0e9353faf786a9d185fa71aa7d62fd97', 12)
Exception:
RuntimeError('shuffle_unpack failed during shuffle e7feb52c45a0263a11c40c62f75ac01c')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 86, in shuffle_unpack
raise RuntimeError(msg)
Task:
('shuffle-p2p-0e9353faf786a9d185fa71aa7d62fd97', 15)
Exception:
RuntimeError('shuffle_unpack failed during shuffle e7feb52c45a0263a11c40c62f75ac01c')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 86, in shuffle_unpack
raise RuntimeError(msg)
Task:
('hash-join-9abf3d2ab874780609827bc89b7c0e24', 0)
Exception:
RuntimeError('Worker tcp://172.21.159.207:37345 left during active shuffle da7a7376729bf812941e1d68ce7a45a7')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 902, in get_output_partition
shuffle = self.get_shuffle_run(shuffle_id, run_id)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 873, in get_shuffle_run
return sync(
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 678, in _get_shuffle_run
shuffle = await self._refresh_shuffle(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 780, in _refresh_shuffle
raise RuntimeError(result["message"])
Task:
('hash-join-3b999bbd22940602da883e391a8c2645', 2)
Exception:
ArrowInvalid('Schema at index 9 was different: \nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: string\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]\nvs\nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: null\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 903, in get_output_partition
return sync(self.worker.loop, shuffle.get_output_partition, output_partition)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 529, in get_output_partition
out = await self.offload(_)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 132, in offload
return await asyncio.get_running_loop().run_in_executor(
File "/opt/conda/lib/python3.10/concurrent/futures/thread.py", line 58, in run
result = self.fn(*self.args, **self.kwargs)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 526, in _
df = convert_partition(data)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_arrow.py", line 57, in convert_partition
return pa.concat_tables(shards)
File "pyarrow/table.pxi", line 5224, in pyarrow.lib.concat_tables
File "pyarrow/error.pxi", line 144, in pyarrow.lib.pyarrow_internal_check_status
File "pyarrow/error.pxi", line 100, in pyarrow.lib.check_status
Task:
('shuffle-p2p-48731fec70b7b0db7cbb86cea0b543e2', 5)
Exception:
RuntimeError('shuffle_unpack failed during shuffle 5ca3decd5a3715fd770d04677df7ba53')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 86, in shuffle_unpack
raise RuntimeError(msg)
Task:
('shuffle-p2p-93ca40036bea0a8853482b832d82821c', 13)
Exception:
RuntimeError('shuffle_unpack failed during shuffle d1ccb1d6ec58027506e3b08ace4baa58')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 86, in shuffle_unpack
raise RuntimeError(msg)
Task:
('shuffle-p2p-f4e54e0fd7edba604e5a9005f0d9304c', 7)
Exception:
RuntimeError('shuffle_unpack failed during shuffle 8048524b8e30f8d0e3c82f8f5d42075c')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 86, in shuffle_unpack
raise RuntimeError(msg)
Task:
('shuffle-transfer-1d083ab93f2404486d7c05fa123be9c9', 7)
Exception:
RuntimeError('shuffle_transfer failed during shuffle 1d083ab93f2404486d7c05fa123be9c9')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 71, in shuffle_transfer
raise RuntimeError(msg)
Task:
shuffle-barrier-f2a8b65b72a7265a5581bfdfb6fe7102
Exception:
RuntimeError('shuffle_barrier failed during shuffle f2a8b65b72a7265a5581bfdfb6fe7102')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 97, in shuffle_barrier
raise RuntimeError(msg)
Task:
('hash-join-a753bbfa960c6db6f21d040b3cf4c994', 4)
Exception:
RuntimeError('Worker tcp://172.21.159.253:41255 left during active shuffle 9a7a3752627bdec5544488ef31d56567')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 902, in get_output_partition
shuffle = self.get_shuffle_run(shuffle_id, run_id)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 873, in get_shuffle_run
return sync(
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 678, in _get_shuffle_run
shuffle = await self._refresh_shuffle(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 780, in _refresh_shuffle
raise RuntimeError(result["message"])
Task:
('shuffle-transfer-373f6c773bfce2becdf1b36cb967f971', 15)
Exception:
RuntimeError('shuffle_transfer failed during shuffle 373f6c773bfce2becdf1b36cb967f971')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 71, in shuffle_transfer
raise RuntimeError(msg)
Task:
shuffle-barrier-1fee22425ad67e23cfadefcdf29a5b52
Exception:
RuntimeError('shuffle_barrier failed during shuffle 1fee22425ad67e23cfadefcdf29a5b52')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 97, in shuffle_barrier
raise RuntimeError(msg)
Task:
('hash-join-73f8548ce5d6c0b50806fdd8dd8ca0a6', 6)
Exception:
RuntimeError('Worker tcp://172.21.25.54:42787 left during active shuffle 3e7274a00bd2078c2f60e47ad3bb2195')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 902, in get_output_partition
shuffle = self.get_shuffle_run(shuffle_id, run_id)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 873, in get_shuffle_run
return sync(
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 678, in _get_shuffle_run
shuffle = await self._refresh_shuffle(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 780, in _refresh_shuffle
raise RuntimeError(result["message"])
Task:
shuffle-barrier-76af57fc7f9a325f94909e7885439392
Exception:
RuntimeError('shuffle_barrier failed during shuffle 76af57fc7f9a325f94909e7885439392')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 97, in shuffle_barrier
raise RuntimeError(msg)
Task:
('hash-join-fdce2ce6f66398c235ed040029cd1481', 1)
Exception:
ArrowInvalid('Schema at index 1 was different: \nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: null\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]\nvs\nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: string\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 903, in get_output_partition
return sync(self.worker.loop, shuffle.get_output_partition, output_partition)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 529, in get_output_partition
out = await self.offload(_)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 132, in offload
return await asyncio.get_running_loop().run_in_executor(
File "/opt/conda/lib/python3.10/concurrent/futures/thread.py", line 58, in run
result = self.fn(*self.args, **self.kwargs)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 526, in _
df = convert_partition(data)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_arrow.py", line 57, in convert_partition
return pa.concat_tables(shards)
File "pyarrow/table.pxi", line 5224, in pyarrow.lib.concat_tables
File "pyarrow/error.pxi", line 144, in pyarrow.lib.pyarrow_internal_check_status
File "pyarrow/error.pxi", line 100, in pyarrow.lib.check_status
Task:
('hash-join-6cc4259a980e9f37ba27009597bd76da', 12)
Exception:
ArrowInvalid('Schema at index 1 was different: \nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: null\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]\nvs\nschema_version: int64\ndetector_id: int64\nsource_id: int64\ntimestamp: int64\ngdt: int64\nsrt: int64\nheading: double\nspeed: double\nevent_type: string\nexterior_lights: string\nlongitude: double\nlatitude: double\naltitude: int64\ngeometry_type: string\nis_day: bool\nhashed_source_and_date: string\nindex: int64\nheading_diff_prev_and_cur: double\nheading_diff_cur_and_next: double\n__hash_partition: uint64\ndate: timestamp[ns]')
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_merge.py", line 173, in merge_unpack
left = ext.get_output_partition(
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 903, in get_output_partition
return sync(self.worker.loop, shuffle.get_output_partition, output_partition)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 416, in sync
raise exc.with_traceback(tb)
File "/opt/conda/lib/python3.10/site-packages/distributed/utils.py", line 389, in f
result = yield future
File "/opt/conda/lib/python3.10/site-packages/tornado/gen.py", line 769, in run
value = future.result()
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 529, in get_output_partition
out = await self.offload(_)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 132, in offload
return await asyncio.get_running_loop().run_in_executor(
File "/opt/conda/lib/python3.10/concurrent/futures/thread.py", line 58, in run
result = self.fn(*self.args, **self.kwargs)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_worker_extension.py", line 526, in _
df = convert_partition(data)
File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_arrow.py", line 57, in convert_partition
return pa.concat_tables(shards)
File "pyarrow/table.pxi", line 5224, in pyarrow.lib.concat_tables
File "pyarrow/error.pxi", line 144, in pyarrow.lib.pyarrow_internal_check_status
File "pyarrow/error.pxi", line 100, in pyarrow.lib.check_status