# Best way to repartition heavily-filtered matrix tables?

**URL:** <https://discuss.hail.is/t/best-way-to-repartition-heavily-filtered-matrix-tables/2140>\
**Category:** Hail Query & hailctl\
**Created:** [July 19, 2021, 1:33am UTC](https://discuss.hail.is/t/best-way-to-repartition-heavily-filtered-matrix-tables/2140 "2021-07-19T01:33:47Z")\
**Posts on this page:** 11\
**Page:** 1

<div class="post-metadata">

**Author:** ![KatalinaBobowik](https://yyz2.discourse-cdn.com/flex036/user_avatar/discuss.hail.is/katalinabobowik/32/588_2.png) [@KatalinaBobowik](https://discuss.hail.is/u/KatalinaBobowik)\
**Post date:** [July 19, 2021, 1:33am UTC](https://discuss.hail.is/t/best-way-to-repartition-heavily-filtered-matrix-tables/2140/1 "2021-07-19T01:33:47Z")

</div>

Hi! I recently posted a [question on Zulip chat](https://hail.zulipchat.com/#narrow/stream/123010-Hail-0.2E2.20support/topic/How.20to.20choose.20number.20of.20partitions.20when.20repartitioning) on how best to repartition heavily-filtered data, however I thought I’d repost here to get any additional thoughts/discussion on what the best approach is to repartitioning.

Prior to the discussion from the link above, I was using the `repartition` function as per the [hail documentation example](https://hail.is/docs/0.2/hail.MatrixTable.html#hail.MatrixTable.repartition) (`dataset_result = dataset.repartition(1000)`), but with shuffle=FALSE (you can find that script [here](https://github.com/populationgenomics/ancestry/blob/main/scripts/hail_batch/hgdp1kg_tobwgs_densify/hgdp_1kg_tob_wgs_densify.py#L41)). I later found that using shuffle=False can cause a loss of parallelism and that repartitioning data can have some unpredictable behaviour. To avoid this, and to have more consistent behaviour from `repartition`, I was recommended to use the following code:

```auto
    mt = mt.semi_join_rows(mt2.rows())
    mt_path = f'{output}/mt_filtered.mt'
    tmp_path = mt_path + '.tmp'
    mt = mt.checkpoint(tmp_path)
    hl.read_matrix_table(tmp_path, _n_partitions=1000).write(mt_path)
    hl.current_backend().fs.rmtree(tmp_path)

```

Unfortunately, using this method, [my script](https://github.com/populationgenomics/ancestry/blob/a45d10e066e92649b9633aa868bddff71969cf0f/scripts/hail_batch/variant_selection/hgdp_1kg_tob_wgs_variant_selection.py) never seemed to finish after allocating 12 hours of run time (using 20 preemptible workers), and failed with the following error:  
`WARNING: Job terminated, but output did not finish streaming`

I tried changing the workers from preemptible to non-preemtible workers, and also played around with the number of partitions (changing from 100 to 1000). When I omit the repartitioning step, the entire script finishes successfully and relatively quickly (less than an hour) using 20 preemptible workers.

I’m wondering if there is another way of repartitioning heavily-filtered files, or perhaps a better way of avoiding many tiny partitions?

Thanks!

---

<div class="post-metadata">

**Author:** ![tpoterba](https://yyz2.discourse-cdn.com/flex036/user_avatar/discuss.hail.is/tpoterba/32/61_2.png) [@tpoterba](https://discuss.hail.is/u/tpoterba)\
**Post date:** [July 20, 2021, 5:55pm UTC](https://discuss.hail.is/t/best-way-to-repartition-heavily-filtered-matrix-tables/2140/2 "2021-07-20T17:55:51Z")

</div>

This is the strategy I would recommend. What’s upstream of the semi\_join\_rows? That this is taking a long time is surprising to me.

---

<div class="post-metadata">

**Author:** ![KatalinaBobowik](https://yyz2.discourse-cdn.com/flex036/user_avatar/discuss.hail.is/katalinabobowik/32/588_2.png) [@KatalinaBobowik](https://discuss.hail.is/u/KatalinaBobowik)\
**Post date:** [July 21, 2021, 4:05am UTC](https://discuss.hail.is/t/best-way-to-repartition-heavily-filtered-matrix-tables/2140/3 "2021-07-21T04:05:45Z")

</div>

Hi Tim, the [script](https://github.com/populationgenomics/ancestry/blob/a45d10e066e92649b9633aa868bddff71969cf0f/scripts/hail_batch/variant_selection/hgdp_1kg_tob_wgs_variant_selection.py) itself actually has quite a lot of costly functions going on before the repartitioning, including densifying the dataset and ld-pruning. However, in the log file these operations all finish, leading me to believe it’s the repartitioning step. The job itself also completes relatively quickly when I run it without repartitioning.

---

<div class="post-metadata">

**Author:** ![tpoterba](https://yyz2.discourse-cdn.com/flex036/user_avatar/discuss.hail.is/tpoterba/32/61_2.png) [@tpoterba](https://discuss.hail.is/u/tpoterba)\
**Post date:** [July 21, 2021, 12:12pm UTC](https://discuss.hail.is/t/best-way-to-repartition-heavily-filtered-matrix-tables/2140/4 "2021-07-21T12:12:49Z")

</div>

This is super weird, I would expect the read(\_n\_partitions=…).write to be much faster than the rest of this query.

is it correct that the difference between the two scripts is the following:

### repartition version

```python
    mt_path = f'{output}/tob_wgs_hgdp_1kg_filtered_variants.mt'
    tmp_path = mt_path + '.tmp'
    hgdp1kg_tobwgs_joined = hgdp1kg_tobwgs_joined.checkpoint(tmp_path)
    hl.read_matrix_table(tmp_path, _n_partitions=100).write(mt_path)
    hl.current_backend().fs.rmtree(tmp_path)

```

### no repartition

```python
    hgdp1kg_tobwgs_joined.write(mt_path)

```

---

<div class="post-metadata">

**Author:** ![KatalinaBobowik](https://yyz2.discourse-cdn.com/flex036/user_avatar/discuss.hail.is/katalinabobowik/32/588_2.png) [@KatalinaBobowik](https://discuss.hail.is/u/KatalinaBobowik)\
**Post date:** [July 22, 2021, 1:13am UTC](https://discuss.hail.is/t/best-way-to-repartition-heavily-filtered-matrix-tables/2140/5 "2021-07-22T01:13:01Z")

</div>

Yes, that’s exactly right. The only difference between the successful and unsuccessful scripts is the code above. I played around with a few different ways of making the script more efficient (such as reducing the total number of variants for LD pruning), but ultimately just removing the repartitioning step (as in your `no repartition` example above) seemed to solve the issue. I know you mentioned earlier that repartitioning can be unpredictable at times, particularly when shuffle=False. I’m curious if there are any alternative methods to reducing the total number of partitions, as this will greatly decrease the computational overhead downstream.

---

<div class="post-metadata">

**Author:** ![tpoterba](https://yyz2.discourse-cdn.com/flex036/user_avatar/discuss.hail.is/tpoterba/32/61_2.png) [@tpoterba](https://discuss.hail.is/u/tpoterba)\
**Post date:** [July 23, 2021, 1:09pm UTC](https://discuss.hail.is/t/best-way-to-repartition-heavily-filtered-matrix-tables/2140/6 "2021-07-23T13:09:07Z")

</div>

OK, I’m totally stumped. This is pretty much the most efficient way to repartition.

If you have the log file for one of these long-running tasks, I would love to look at that – it should tell us where the pipeline is stalled.

---

<div class="post-metadata">

**Author:** ![tpoterba](https://yyz2.discourse-cdn.com/flex036/user_avatar/discuss.hail.is/tpoterba/32/61_2.png) [@tpoterba](https://discuss.hail.is/u/tpoterba)\
**Post date:** [July 23, 2021, 1:09pm UTC](https://discuss.hail.is/t/best-way-to-repartition-heavily-filtered-matrix-tables/2140/7 "2021-07-23T13:09:27Z")

</div>

(it could be the `rmtree` in fact)

---

<div class="post-metadata">

**Author:** ![KatalinaBobowik](https://yyz2.discourse-cdn.com/flex036/user_avatar/discuss.hail.is/katalinabobowik/32/588_2.png) [@KatalinaBobowik](https://discuss.hail.is/u/KatalinaBobowik)\
**Post date:** [July 26, 2021, 7:51am UTC](https://discuss.hail.is/t/best-way-to-repartition-heavily-filtered-matrix-tables/2140/8 "2021-07-26T07:51:03Z")

</div>

No problem, I’m attaching two log files here: one [with using rmtree](https://drive.google.com/file/d/1tWXhHrDnaDoG1rFlv_UXQU9bgRU8kTb-/view?usp=sharing) with 100 partitions, and one [without using rmtree](https://drive.google.com/file/d/1FILVYk0sfDwgot0dHY4jMhOaC9D7I7zj/view?usp=sharing) with 1000 partitions.

The script for the first log file (with rmtree) can be found [here](https://github.com/populationgenomics/ancestry/blob/3d428924f1809dd98086ee35ffee24190dc15739/scripts/hail_batch/variant_selection/hgdp_1kg_tob_wgs_variant_selection.py), and the second one (without rmtree) can be found [here](https://github.com/populationgenomics/ancestry/blob/fc942c0fb249d4da7ac52db21ffaed70fcb12e62/scripts/hail_batch/variant_selection/hgdp_1kg_tob_wgs_variant_selection.py).

---

<div class="post-metadata">

**Author:** ![tpoterba](https://yyz2.discourse-cdn.com/flex036/user_avatar/discuss.hail.is/tpoterba/32/61_2.png) [@tpoterba](https://discuss.hail.is/u/tpoterba)\
**Post date:** [July 26, 2021, 12:28pm UTC](https://discuss.hail.is/t/best-way-to-repartition-heavily-filtered-matrix-tables/2140/9 "2021-07-26T12:28:48Z")

</div>

Hi Katalina,  
These actually aren’t the logs I need – Hail writes a log file to the current working directory while it runs, and this is the one that has more detailed information (with timestamps) about the compilation and execution of each query. On Dataproc, these logs are written to files in `/home/hail` on the driver VM. They can be somewhat large (\>100MB) for queries with a lot of Spark logging (which scales with the # of partitions), so feel free to send by email if that’s easier.

For the long-running job, just SSH in and copy out the log when it hits the point that it stops making progress – we shouldn’t need hours of logging about that.

---

<div class="post-metadata">

**Author:** ![KatalinaBobowik](https://yyz2.discourse-cdn.com/flex036/user_avatar/discuss.hail.is/katalinabobowik/32/588_2.png) [@KatalinaBobowik](https://discuss.hail.is/u/KatalinaBobowik)\
**Post date:** [July 27, 2021, 6:35am UTC](https://discuss.hail.is/t/best-way-to-repartition-heavily-filtered-matrix-tables/2140/10 "2021-07-27T06:35:31Z")

</div>

Hi Tim, here’s the link to the log file with the more detailed information: [Google Drive: Sign-in](https://drive.google.com/file/d/12I3jsAgszM-HYsDoxpfgh4AW1c4aAjn1/view?usp=sharing)

---

<div class="post-metadata">

**Author:** ![KatalinaBobowik](https://yyz2.discourse-cdn.com/flex036/user_avatar/discuss.hail.is/katalinabobowik/32/588_2.png) [@KatalinaBobowik](https://discuss.hail.is/u/KatalinaBobowik)\
**Post date:** [August 24, 2021, 4:04am UTC](https://discuss.hail.is/t/best-way-to-repartition-heavily-filtered-matrix-tables/2140/11 "2021-08-24T04:04:59Z")

</div>

Hi @tpoterba, just a friendly ping- Did you manage to look into the log files to see where the issue might be occurring? Let me know if you need to me to attach anything else!
