()
| 120 | |
| 121 | @pytest.mark.slow() |
| 122 | def test_refcounting_futures(): |
| 123 | pd = pytest.importorskip("pandas") |
| 124 | dd = pytest.importorskip("dask.dataframe") |
| 125 | distributed = pytest.importorskip("distributed") |
| 126 | |
| 127 | # See https://github.com/dask/distributed/issues/9041 |
| 128 | # Didn't reproduce with any of our fixtures |
| 129 | with distributed.Client( |
| 130 | n_workers=2, worker_class=distributed.Worker, dashboard_address=":0" |
| 131 | ) as client: |
| 132 | |
| 133 | def gen(i): |
| 134 | return pd.DataFrame({"A": [i]}, index=[i]) |
| 135 | |
| 136 | futures = [client.submit(gen, i) for i in range(3)] |
| 137 | |
| 138 | meta = gen(0)[:0] |
| 139 | df = dd.from_delayed(futures, meta) |
| 140 | df.compute() |
| 141 | |
| 142 | del futures |
| 143 | |
| 144 | df.compute() |
| 145 | |
| 146 | |
| 147 | class FooExpr(Expr): |
nothing calls this directly
no test coverage detected