Exceptions

Task: ('from_pandas-5e71892f0e34621cf2d740f0c2b1f96f', 8)

Worker(s):

Exception:

Traceback
 

Task: ('from_pandas-5e71892f0e34621cf2d740f0c2b1f96f', 14)

Worker(s):

Exception:

Traceback
 

Task: ('from_pandas-6d59010172bbd4429cbe00e2a674eb65', 9)

Worker(s):

Exception:

Traceback
 

Task: ('from_pandas-6d59010172bbd4429cbe00e2a674eb65', 12)

Worker(s):

Exception:

Traceback
 

Task: ('from_pandas-6d59010172bbd4429cbe00e2a674eb65', 15)

Worker(s):

Exception:

Traceback
 

Task: ('from_pandas-5e71892f0e34621cf2d740f0c2b1f96f', 15)

Worker(s):

Exception:

Traceback
 

Task: ('hash-join-4b4b25fcaf703c1e99e652bac3539b7f', 3)

Worker(s):

Exception: RuntimeError('Worker tcp://172.21.25.29:34733 left during active shuffle 3d4472ea35c0f53629acee4ad7bde540')

Traceback
  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)

Worker(s):

Exception: RuntimeError('shuffle_unpack failed during shuffle 416a4a8a9303bd1b6a928557fd2c4409')

Traceback
  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)

Worker(s):

Exception: RuntimeError('shuffle_unpack failed during shuffle 416a4a8a9303bd1b6a928557fd2c4409')

Traceback
  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)

Worker(s):

Exception: RuntimeError('shuffle_transfer failed during shuffle 2f4fca06bc9c5475b26c12cf93722f4f')

Traceback
  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: ('from_pandas-5e71892f0e34621cf2d740f0c2b1f96f', 13)

Worker(s):

Exception:

Traceback
 

Task: run-13177c827a5279df9faaca249711e262

Worker(s):

Exception: KilledWorker("('from_pandas-7257ad7e7e1317b17a9f46f3bd9f975f', 10)", <WorkerState 'tcp://172.21.159.246:38983', status: closed, memory: 0, processing: 0>, 3)

Traceback
  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: ('from_pandas-7257ad7e7e1317b17a9f46f3bd9f975f', 10)

Worker(s):

Exception:

Traceback
 

Task: ('from_pandas-5e71892f0e34621cf2d740f0c2b1f96f', 2)

Worker(s):

Exception:

Traceback
 

Task: ('assign-f675a204035f907c92325171b70e20d1', 1)

Worker(s):

Exception: GEOSException('IllegalArgumentException: point array must contain 0 or >1 elements\n')

Traceback
  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: ('from_pandas-fece414c141f919c42a3acf6b286a9c7', 4)

Worker(s):

Exception:

Traceback
 

Task: ('hash-join-transfer-05b2eea45734abca7b550c5b98f8f3a6', 13)

Worker(s):

Exception: RuntimeError('shuffle_transfer failed during shuffle 05b2eea45734abca7b550c5b98f8f3a6')

Traceback
  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)

Worker(s):

Exception: RuntimeError('shuffle_transfer failed during shuffle eec3c2cf83943d5b1a0c93ed88834f19')

Traceback
  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)

Worker(s):

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]')

Traceback
  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)

Worker(s):

Exception: RuntimeError('shuffle_unpack failed during shuffle 8f860c6fe1f949e8bb18fe310bb37739')

Traceback
  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)

Worker(s):

Exception: RuntimeError('shuffle_unpack failed during shuffle 8f860c6fe1f949e8bb18fe310bb37739')

Traceback
  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)

Worker(s):

Exception: RuntimeError('shuffle_unpack failed during shuffle 080fc1af058f130722ca3ee502f17ab5')

Traceback
  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)

Worker(s):

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]')

Traceback
  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)

Worker(s):

Exception: ValueError('Found array with 1 sample(s) (shape=(1, 1)) while a minimum of 2 is required by AgglomerativeClustering.')

Traceback
  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: run-fbad1848d1ff3d5a571bc283b409ba5a

Worker(s):

Exception:

Traceback
 

Task: ('sjoin-f7e724a9485d517e60cbe0de77064215', 66)

Worker(s):

Exception:

Traceback
 

Task: ('hash-join-0878ad0b59064ab99a6235ac6c7270e1', 15)

Worker(s):

Exception:

Traceback
 

Task: ('hash-join-0878ad0b59064ab99a6235ac6c7270e1', 12)

Worker(s):

Exception:

Traceback
 

Task: ('hash-join-0878ad0b59064ab99a6235ac6c7270e1', 5)

Worker(s):

Exception:

Traceback
 

Task: ('hash-join-0ad9eedf603cf4e330cfc753a28d9c23', 9)

Worker(s):

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]')

Traceback
  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)

Worker(s):

Exception: RuntimeError('Worker tcp://172.21.151.150:38387 left during active shuffle 7fe06a3da8052015aabaf82be4107670')

Traceback
  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: ('sjoin-fa0424d783adbfed9fca4e1397719f19', 165)

Worker(s):

Exception:

Traceback
 

Task: ('sjoin-c1d04ffc4dd32b8081e833d5ca345962', 250)

Worker(s):

Exception:

Traceback
 

Task: ('shuffle-p2p-ee7354768e111c1f5f98db89d91d0ab5', 7)

Worker(s):

Exception: RuntimeError('shuffle_unpack failed during shuffle 14b747bcff3a62b3bed39aaa4dfdaa24')

Traceback
  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)

Worker(s):

Exception: RuntimeError('shuffle_unpack failed during shuffle 14b747bcff3a62b3bed39aaa4dfdaa24')

Traceback
  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)

Worker(s):

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]')

Traceback
  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)

Worker(s):

Exception: GEOSException('IllegalArgumentException: point array must contain 0 or >1 elements\n')

Traceback
  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: run-36cf509e3c2db9a20146cda43b7b1838

Worker(s):

Exception:

Traceback
 

Task: ('hash-join-a80329184eccb02b1f490313add46065', 5)

Worker(s):

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]')

Traceback
  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: run-20995ff95188017d71b39bbd7f700e4a

Worker(s):

Exception:

Traceback
 

Task: ('from_pandas-6b24fc4d4f63d5531d6c1e7931483f3a', 12)

Worker(s):

Exception:

Traceback
 

Task: ('hash-join-24a31f1dadf8bf09023eb1d3e9e24c90', 0)

Worker(s):

Exception: RuntimeError('Worker tcp://172.21.151.150:44549 left during active shuffle 40a265a636f672d90a667506167eb84d')

Traceback
  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)

Worker(s):

Exception: RuntimeError('Worker tcp://172.21.159.246:36409 left during active shuffle 8135f54f9dce30185f4d3fef2ada6bdd')

Traceback
  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)

Worker(s):

Exception: RuntimeError('shuffle_transfer failed during shuffle 2c7fd55cad0e13fedbad3273e327127d')

Traceback
  File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 71, in shuffle_transfer
    raise RuntimeError(msg)
 

Task: shuffle-barrier-4525ea1da2d0e9579776be9fe4301954

Worker(s):

Exception: RuntimeError('shuffle_barrier failed during shuffle 4525ea1da2d0e9579776be9fe4301954')

Traceback
  File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 97, in shuffle_barrier
    raise RuntimeError(msg)
 

Task: ('assign-3361cddacb9ac6054e95a968f5d5153a', 6)

Worker(s):

Exception: GEOSException('IllegalArgumentException: point array must contain 0 or >1 elements\n')

Traceback
  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: run-552331c6b8985d8735303d3d126bc3d1

Worker(s):

Exception:

Traceback
 

Task: ('hash-join-c89798abd68c8b783c220d4c4b8fe23b', 15)

Worker(s):

Exception: RuntimeError('Worker tcp://172.21.25.54:39899 left during active shuffle 64ad8ad5105af031ddd958c4a37785f7')

Traceback
  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)

Worker(s):

Exception: RuntimeError('shuffle_unpack failed during shuffle e3ef17c11773502e54c27ab45b7c658a')

Traceback
  File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 86, in shuffle_unpack
    raise RuntimeError(msg)
 

Task: run-30a11945490277dcdd22527e086895f2

Worker(s):

Exception:

Traceback
 

Task: ('hash-join-569f92d5b293d53137f7a28d5370cbbf', 8)

Worker(s):

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]')

Traceback
  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)

Worker(s):

Exception: RuntimeError('shuffle_unpack failed during shuffle ce71486c0e31fb151f9ca45582783a7f')

Traceback
  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)

Worker(s):

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]')

Traceback
  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)

Worker(s):

Exception: RuntimeError('Worker tcp://172.21.159.253:40571 left during active shuffle 9a400d2158ab3d046210d22e656072bc')

Traceback
  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: run-da2f0b24d3ded2e858125662d32ce4ea

Worker(s):

Exception:

Traceback
 

Task: ('hash-join-transfer-ee6fa7f2d2572fa5840c05b4e74410bd', 5)

Worker(s):

Exception: RuntimeError('shuffle_transfer failed during shuffle ee6fa7f2d2572fa5840c05b4e74410bd')

Traceback
  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-9df21a6cd9ff2716f57bf2d1d222d80f', 5)

Worker(s):

Exception:

Traceback
 

Task: ('hash-join-9df21a6cd9ff2716f57bf2d1d222d80f', 8)

Worker(s):

Exception:

Traceback
 

Task: ('hash-join-9df21a6cd9ff2716f57bf2d1d222d80f', 0)

Worker(s):

Exception:

Traceback
 

Task: ('sjoin-d8d3dff67c9ea1218d13b472f7f06ded', 161)

Worker(s):

Exception:

Traceback
 

Task: shuffle-barrier-011f1efefe4c945109110a4b577cbf4d

Worker(s):

Exception: RuntimeError('shuffle_barrier failed during shuffle 011f1efefe4c945109110a4b577cbf4d')

Traceback
  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)

Worker(s):

Exception: RuntimeError('shuffle_transfer failed during shuffle b09e00533ac31cc9c2c310d2c2305e75')

Traceback
  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)

Worker(s):

Exception: RuntimeError('Worker tcp://172.21.159.222:35865 left during active shuffle 97cc0e65355bd18db884b23ece825ae9')

Traceback
  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)

Worker(s):

Exception: RuntimeError('Worker tcp://172.21.159.207:43461 left during active shuffle e90d5c7caa33f2f3560741d3a1f64fab')

Traceback
  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

Worker(s):

Exception: RuntimeError('shuffle_barrier failed during shuffle a079868710da4a8af2c4f117f9c314dc')

Traceback
  File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 97, in shuffle_barrier
    raise RuntimeError(msg)
 

Task: run-93cb3608d73bb3603f4fbb7dabb40b66

Worker(s):

Exception:

Traceback
 

Task: ('shuffle-p2p-0e9353faf786a9d185fa71aa7d62fd97', 12)

Worker(s):

Exception: RuntimeError('shuffle_unpack failed during shuffle e7feb52c45a0263a11c40c62f75ac01c')

Traceback
  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)

Worker(s):

Exception: RuntimeError('shuffle_unpack failed during shuffle e7feb52c45a0263a11c40c62f75ac01c')

Traceback
  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)

Worker(s):

Exception: RuntimeError('Worker tcp://172.21.159.207:37345 left during active shuffle da7a7376729bf812941e1d68ce7a45a7')

Traceback
  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)

Worker(s):

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]')

Traceback
  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)

Worker(s):

Exception: RuntimeError('shuffle_unpack failed during shuffle 5ca3decd5a3715fd770d04677df7ba53')

Traceback
  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)

Worker(s):

Exception: RuntimeError('shuffle_unpack failed during shuffle d1ccb1d6ec58027506e3b08ace4baa58')

Traceback
  File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 86, in shuffle_unpack
    raise RuntimeError(msg)
 

Task: ('from_pandas-6eff0ac2d8a5523414337928f5fa4785', 14)

Worker(s):

Exception:

Traceback
 

Task: ('hash-join-ef6b4f39704282dd70aa09b6f45ae113', 1)

Worker(s):

Exception:

Traceback
 

Task: run-b8ea6ebc5140e5cfabc928500fdd6dad

Worker(s):

Exception:

Traceback
 

Task: ('shuffle-p2p-f4e54e0fd7edba604e5a9005f0d9304c', 7)

Worker(s):

Exception: RuntimeError('shuffle_unpack failed during shuffle 8048524b8e30f8d0e3c82f8f5d42075c')

Traceback
  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)

Worker(s):

Exception: RuntimeError('shuffle_transfer failed during shuffle 1d083ab93f2404486d7c05fa123be9c9')

Traceback
  File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 71, in shuffle_transfer
    raise RuntimeError(msg)
 

Task: ('from_pandas-d744eab5f8a0b5d6ecd7c4055734eea7', 3)

Worker(s):

Exception:

Traceback
 

Task: ('from_pandas-d744eab5f8a0b5d6ecd7c4055734eea7', 12)

Worker(s):

Exception:

Traceback
 

Task: ('from_pandas-d744eab5f8a0b5d6ecd7c4055734eea7', 15)

Worker(s):

Exception:

Traceback
 

Task: ('from_pandas-d744eab5f8a0b5d6ecd7c4055734eea7', 7)

Worker(s):

Exception:

Traceback
 

Task: ('from_pandas-f2c9f00a2fdeb85b20c2f0c6ef740c03', 7)

Worker(s):

Exception:

Traceback
 

Task: ('from_pandas-d5a958f3021cdbf5030098754668ac59', 0)

Worker(s):

Exception:

Traceback
 

Task: ('from_pandas-3be5f9432f2fd15c1ed804b25b95b95d', 9)

Worker(s):

Exception:

Traceback
 

Task: run-ddecefd60677110ce87a981f7319108c

Worker(s):

Exception:

Traceback
 

Task: run-a2dbd316f5ed53569c30dea548c0d572

Worker(s):

Exception:

Traceback
 

Task: shuffle-barrier-f2a8b65b72a7265a5581bfdfb6fe7102

Worker(s):

Exception: RuntimeError('shuffle_barrier failed during shuffle f2a8b65b72a7265a5581bfdfb6fe7102')

Traceback
  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)

Worker(s):

Exception: RuntimeError('Worker tcp://172.21.159.253:41255 left during active shuffle 9a7a3752627bdec5544488ef31d56567')

Traceback
  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)

Worker(s):

Exception: RuntimeError('shuffle_transfer failed during shuffle 373f6c773bfce2becdf1b36cb967f971')

Traceback
  File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 71, in shuffle_transfer
    raise RuntimeError(msg)
 

Task: shuffle-barrier-1fee22425ad67e23cfadefcdf29a5b52

Worker(s):

Exception: RuntimeError('shuffle_barrier failed during shuffle 1fee22425ad67e23cfadefcdf29a5b52')

Traceback
  File "/opt/conda/lib/python3.10/site-packages/distributed/shuffle/_shuffle.py", line 97, in shuffle_barrier
    raise RuntimeError(msg)
 

Task: run-b1e3deb1a63553bc3d23ba43b326c59c

Worker(s):

Exception:

Traceback
 

Task: ('sjoin-74631e305fd13ebc11330a3a99e262ce', 215)

Worker(s):

Exception:

Traceback
 

Task: ('sjoin-74631e305fd13ebc11330a3a99e262ce', 219)

Worker(s):

Exception:

Traceback
 

Task: ('hash-join-73f8548ce5d6c0b50806fdd8dd8ca0a6', 6)

Worker(s):

Exception: RuntimeError('Worker tcp://172.21.25.54:42787 left during active shuffle 3e7274a00bd2078c2f60e47ad3bb2195')

Traceback
  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

Worker(s):

Exception: RuntimeError('shuffle_barrier failed during shuffle 76af57fc7f9a325f94909e7885439392')

Traceback
  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)

Worker(s):

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]')

Traceback
  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-8f31526927863ad5f130de5e921cd6b6', 1)

Worker(s):

Exception:

Traceback
 

Task: ('hash-join-8f31526927863ad5f130de5e921cd6b6', 14)

Worker(s):

Exception:

Traceback
 

Task: ('from_pandas-c095576b7f753fffed310ea98022566d', 9)

Worker(s):

Exception:

Traceback
 

Task: ('hash-join-6cc4259a980e9f37ba27009597bd76da', 12)

Worker(s):

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]')

Traceback
  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