Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
90 commits
Select commit Hold shift + click to select a range
7604abf
Add test illustrating import issue
jrbourbeau Mar 12, 2021
0939cf4
move layer materialization for shuffle
rjzamora Mar 12, 2021
8947089
add missing layers.py file and revise org a bit
rjzamora Mar 12, 2021
05283d7
Merge remote-tracking branch 'jrbourbeau/scheduler-import-test' into …
rjzamora Mar 12, 2021
579fa6b
move layers.py
rjzamora Mar 12, 2021
4e4fee1
remove xfail
rjzamora Mar 15, 2021
948e670
use import rather than pickle for functions
rjzamora Mar 15, 2021
42da7e0
fix importlib error and use dict
rjzamora Mar 15, 2021
3a2ebb4
roll back test changes from 7374 (let failures be resolved in that PR)
rjzamora Mar 15, 2021
958322a
remove debug print
rjzamora Mar 15, 2021
e1e4e18
Merge remote-tracking branch 'upstream/main' into shuffle-avoid-pd-im…
rjzamora Mar 15, 2021
9a3d3f7
move Shuffle layers to layers.py completely
rjzamora Mar 15, 2021
2edba11
moving moving BroadcastJoinLayer to layers.py
rjzamora Mar 15, 2021
f8ed2f9
tweak CallableLazyImport
rjzamora Mar 15, 2021
3ebbef3
move BlockwiseCreateArray into layers.py
rjzamora Mar 15, 2021
9badf2f
import test coverage
rjzamora Mar 15, 2021
e62dba6
incorperate testing idea from 7374
rjzamora Mar 16, 2021
6647084
comment tweaks
rjzamora Mar 16, 2021
79f5044
introduce new test_layers.py module
rjzamora Mar 16, 2021
f3f5a96
update comment in test
rjzamora Mar 16, 2021
d37940c
Update dask/layers.py
rjzamora Mar 16, 2021
bf94c89
Only use CallableLazyImport when the graph is materialized within __d…
rjzamora Mar 17, 2021
6321802
remove obsolete annotation handling
rjzamora Mar 17, 2021
dc54e49
Update dask/layers.py
rjzamora Mar 17, 2021
5b35fba
_construct_graph fix
rjzamora Mar 17, 2021
9b7d5fc
Merge remote-tracking branch 'upstream/main' into shuffle-avoid-pd-im…
rjzamora Mar 17, 2021
3813b61
migrate csv code
rjzamora Mar 17, 2021
be01159
migrate orc changes
rjzamora Mar 17, 2021
0b71c39
basic parquet migration - still need to handle serialization
rjzamora Mar 17, 2021
fb94266
add require_pickle option to DataFrameIOLayer
rjzamora Mar 17, 2021
e7a4733
Merge remote-tracking branch 'upstream/main' into blockwise-io-dataframe
rjzamora Mar 18, 2021
3b81fec
update testing
rjzamora Mar 18, 2021
56d7df5
change 'column culling' language to 'column projection'
rjzamora Mar 22, 2021
9df767c
Merge remote-tracking branch 'upstream/main' into blockwise-io-dataframe
rjzamora Mar 23, 2021
8aa2dd7
use ConcatAxesWrapper to avoid inline task
rjzamora Mar 23, 2021
b7a80b7
add test coverage
rjzamora Mar 23, 2021
e000748
use serialize instead of pickle
rjzamora Mar 24, 2021
e4ec2f2
use simpler SerializedFunction wrapper
rjzamora Mar 27, 2021
372778f
update to rely on distributed#4575
rjzamora Mar 30, 2021
33d1629
serialization experiments
rjzamora Apr 7, 2021
aed2f05
Merge remote-tracking branch 'upstream/main' into blockwise-concatena…
rjzamora Apr 7, 2021
e1d8e1a
align with current status of distributed-4641
rjzamora Apr 8, 2021
b21b331
update/fix comment
rjzamora Apr 8, 2021
7e49ee6
align with latest distributed#4641 state
rjzamora Apr 13, 2021
ba7e872
re-align with James' suggestion
rjzamora Apr 14, 2021
76275a5
Temporarily point to Distributed PR #4641
jrbourbeau Apr 14, 2021
c6522e9
Merge remote-tracking branch 'upstream/main' into blockwise-concatenate
rjzamora Apr 19, 2021
c72e67a
point CI back to distributed main branch
rjzamora Apr 19, 2021
922e8bf
avoid using distributed 2021.03.0 for mindeps test
rjzamora Apr 19, 2021
39e2893
fix initial merge conflicts
rjzamora Apr 19, 2021
0aa95fd
Merge remote-tracking branch 'origin/blockwise-concatenate' into bloc…
rjzamora Apr 19, 2021
cba6acf
start aligning with #7455
rjzamora Apr 19, 2021
f9ede74
begin requiring io_deps elements to inherit from base class
rjzamora Apr 19, 2021
3e82456
more simplification
rjzamora Apr 19, 2021
80c4b54
fix indices init
rjzamora Apr 19, 2021
c9fd246
Add .git suffix
jrbourbeau Apr 19, 2021
87fda33
Add explicit pip
jrbourbeau Apr 19, 2021
636becf
serialization cleanup
rjzamora Apr 20, 2021
6d8260e
address dep-name issue
rjzamora Apr 20, 2021
1465252
strip out unnecessary import logic
rjzamora Apr 20, 2021
efc15f2
use output_blocks in pack
rjzamora Apr 20, 2021
c76fe90
more code review suggestions
rjzamora Apr 21, 2021
1b745d7
use clearer tmp variable name
rjzamora Apr 21, 2021
306ceb8
avoid io_deps copy
rjzamora Apr 21, 2021
6d264fb
dictionary comp for BlockwiseDepDict pack
rjzamora Apr 21, 2021
e772d11
start working on csv problem
rjzamora Apr 21, 2021
a6af2c6
handle csv nested-task problem
rjzamora Apr 21, 2021
1635e0d
further cleanup
rjzamora Apr 21, 2021
0cb2d6b
apply fix
rjzamora Apr 21, 2021
212db07
add produces_tasks
rjzamora Apr 21, 2021
5ff24fe
use place-holder required_indices variable
rjzamora Apr 21, 2021
fe597c2
CreateArrayDeps fix
rjzamora Apr 22, 2021
71d4aaf
remove extra name def
rjzamora Apr 22, 2021
901e9f3
Merge remote-tracking branch 'upstream/main' into blockwise-concatenate
rjzamora Apr 22, 2021
d7514f6
Change test ordering
jrbourbeau Apr 22, 2021
6f60a69
fix dumps_function mistake and account for required_indices being an …
rjzamora Apr 22, 2021
5bfe1b1
Merge remote-tracking branch 'origin/blockwise-concatenate' into bloc…
rjzamora Apr 22, 2021
dcac772
Merge branch 'main' of https://github.com/dask/dask into blockwise-co…
jrbourbeau Apr 23, 2021
f8a6f6d
Merge branch 'blockwise-concatenate' of https://github.com/rjzamora/d…
rjzamora Apr 23, 2021
d550dd1
Merge remote-tracking branch 'upstream/main' into blockwise-io-dataframe
rjzamora Apr 23, 2021
d7cf995
Merge remote-tracking branch 'upstream/main' into blockwise-io-dataframe
rjzamora Apr 26, 2021
b9c4c57
roll back (breaking) large_graph_objects change
rjzamora Apr 26, 2021
365962a
Merge remote-tracking branch 'upstream/main' into blockwise-io-dataframe
rjzamora Apr 26, 2021
d20b510
minor blockwise changes to address code-review
rjzamora Apr 27, 2021
94cf61b
more cleanup
rjzamora Apr 27, 2021
ab4a35f
use for column projection in csv
rjzamora Apr 27, 2021
0dd1cdf
updating some comments
rjzamora Apr 27, 2021
83fafc4
improve documentation
rjzamora Apr 27, 2021
3ced1c8
add project_columns to the functions to make things a bit more explic…
rjzamora Apr 27, 2021
b9210b4
fix empty-column repr error
rjzamora Apr 27, 2021
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
307 changes: 227 additions & 80 deletions dask/blockwise.py

Large diffs are not rendered by default.

127 changes: 72 additions & 55 deletions dask/dataframe/io/csv.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,11 @@
import copy
from collections.abc import Mapping
from io import BytesIO
from warnings import catch_warnings, simplefilter, warn

from ...highlevelgraph import HighLevelGraph
from ...layers import DataFrameIOLayer

try:
import psutil
except ImportError:
Expand Down Expand Up @@ -32,86 +36,82 @@
from ..utils import clear_known_categories


class CSVSubgraph(Mapping):
class CSVFunctionWrapper:
"""
Subgraph for reading CSV files.
CSV Function-Wrapper Class
Reads CSV data from disk to produce a partition (given a key).
"""

def __init__(
self,
name,
reader,
blocks,
is_first,
columns,
colname,
head,
header,
kwargs,
reader,
dtypes,
columns,
enforce,
path,
kwargs,
):
self.name = name
self.full_columns = columns
self.colname = colname
self.head = head
self.header = header
self.reader = reader
self.blocks = blocks
self.is_first = is_first
self.head = head # example pandas DF for metadata
self.header = header # prepend to all blocks
self.kwargs = kwargs
self.dtypes = dtypes
self.columns = columns
self.enforce = enforce
self.colname, self.paths = path or (None, None)
self.kwargs = kwargs
self.columns = None # Used to pass `usecols`

def __getitem__(self, key):
try:
name, i = key
except ValueError:
# too many / few values to unpack
raise KeyError(key) from None
def project_columns(self, columns):
"""Return a new CSVFunctionWrapper object with
a sub-column projection.
"""
if columns == self.columns:
return self
func = copy.deepcopy(self)
func.columns = columns
return func

if name != self.name:
raise KeyError(key)
def __call__(self, part):

if i < 0 or i >= len(self.blocks):
raise KeyError(key)
# Part will be a 3-element tuple
block, path, is_first = part

block = self.blocks[i]
if self.paths is not None:
# Construct `path_info`
if path is not None:
path_info = (
self.colname,
self.paths[i],
path,
sorted(list(self.head[self.colname].cat.categories)),
)
else:
path_info = None

# Deal with arguments that are special
# for the first block of each file
write_header = False
rest_kwargs = self.kwargs.copy()
if not self.is_first[i]:
if self.columns is not None:
if rest_kwargs.get("usecols", None) is None:
rest_kwargs["usecols"] = self.columns
if not is_first:
write_header = True
rest_kwargs.pop("skiprows", None)

return (
pandas_read_text,
# Call `pandas_read_text`
return pandas_read_text(
self.reader,
block,
self.header,
rest_kwargs,
self.dtypes,
self.columns,
self.full_columns,
Comment thread
rjzamora marked this conversation as resolved.
write_header,
self.enforce,
path_info,
)

def __len__(self):
return len(self.blocks)

def __iter__(self):
for i in range(len(self)):
yield (self.name, i)


def pandas_read_text(
reader,
Expand Down Expand Up @@ -325,8 +325,6 @@ def text_blocks_to_pandas(
# Create mask of first blocks from nested block_lists
is_first = tuple(block_mask(block_lists))

name = "read-csv-" + tokenize(reader, columns, enforce, head, blocksize)

if path:
colname, path_converter = path
paths = [b[1].path for b in blocks]
Expand All @@ -344,21 +342,40 @@ def text_blocks_to_pandas(
if len(unknown_categoricals):
head = clear_known_categories(head, cols=unknown_categoricals)

subgraph = CSVSubgraph(
# Define parts
parts = []
colname, paths = path or (None, None)
for i in range(len(blocks)):
parts.append(
[
blocks[i],
paths[i] if paths else None,
is_first[i],
]
)

# Create Blockwise layer
label = "read-csv-"
name = label + tokenize(reader, columns, enforce, head, blocksize)
layer = DataFrameIOLayer(
name,
reader,
blocks,
is_first,
head,
header,
kwargs,
dtypes,
columns,
enforce,
path,
parts,
CSVFunctionWrapper(
columns,
colname,
head,
header,
reader,
dtypes,
enforce,
kwargs,
),
label=label,
produces_tasks=True,
)

return new_dd_object(subgraph, name, head, (None,) * (len(blocks) + 1))
graph = HighLevelGraph({name: layer}, {name: set()})
return new_dd_object(graph, name, head, (None,) * (len(blocks) + 1))


def block_mask(block_lists):
Expand Down
67 changes: 51 additions & 16 deletions dask/dataframe/io/orc.py
Original file line number Diff line number Diff line change
@@ -1,15 +1,49 @@
import copy
from distutils.version import LooseVersion

from fsspec.core import get_fs_token_paths

from ...base import tokenize
from ...highlevelgraph import HighLevelGraph
from ...layers import DataFrameIOLayer
from ...utils import import_required
from ..core import DataFrame
from .utils import _get_pyarrow_dtypes, _meta_from_dtypes

__all__ = ("read_orc",)


class ORCFunctionWrapper:
"""
ORC Function-Wrapper Class
Reads ORC data from disk to produce a partition.
"""

def __init__(self, fs, columns, schema):
self.fs = fs
self.columns = columns
self.schema = schema

def project_columns(self, columns):
"""Return a new ORCFunctionWrapper object with
a sub-column projection.
"""
if columns == self.columns:
return self
func = copy.deepcopy(self)
func.columns = columns
return func

def __call__(self, stripe_info):
path, stripe = stripe_info
return _read_orc_stripe(
self.fs,
path,
stripe,
list(self.schema) if self.columns is None else self.columns,
)


def _read_orc_stripe(fs, path, stripe, columns=None):
"""Pull out specific data from specific part of ORC file"""
orc = import_required("pyarrow.orc", "Please install pyarrow >= 0.9.0")
Expand All @@ -26,7 +60,6 @@ def _read_orc_stripe(fs, path, stripe, columns=None):

def read_orc(path, columns=None, storage_options=None):
"""Read dataframe from ORC file(s)

Parameters
----------
path: str or list(str)
Expand All @@ -36,11 +69,9 @@ def read_orc(path, columns=None, storage_options=None):
Columns to load. If None, loads all.
storage_options: None or dict
Further parameters to pass to the bytes backend.

Returns
-------
Dask.DataFrame (even if there is only one column)

Examples
--------
>>> df = dd.read_orc('https://github.com/apache/orc/raw/'
Expand All @@ -64,32 +95,36 @@ def read_orc(path, columns=None, storage_options=None):
path, mode="rb", storage_options=storage_options
)
schema = None
nstripes_per_file = []
parts = []
for path in paths:
with fs.open(path, "rb") as f:
o = orc.ORCFile(f)
if schema is None:
schema = o.schema
elif schema != o.schema:
raise ValueError("Incompatible schemas while parsing ORC files")
nstripes_per_file.append(o.nstripes)
for stripe in range(o.nstripes):
parts.append((path, stripe))
schema = _get_pyarrow_dtypes(schema, categories=None)
if columns is not None:
ex = set(columns) - set(schema)
if ex:
raise ValueError(
"Requested columns (%s) not in schema (%s)" % (ex, set(schema))
)
else:
columns = list(schema)
meta = _meta_from_dtypes(columns, schema, [], [])

name = "read-orc-" + tokenize(fs_token, path, columns)
dsk = {}
N = 0
for path, n in zip(paths, nstripes_per_file):
for stripe in range(n):
dsk[(name, N)] = (_read_orc_stripe, fs, path, stripe, columns)
N += 1
# Create Blockwise layer
label = "read-orc-"
output_name = label + tokenize(fs_token, path, columns)
layer = DataFrameIOLayer(
output_name,
columns,
parts,
ORCFunctionWrapper(fs, columns, schema),
label=label,
)

return DataFrame(dsk, name, meta, [None] * (len(dsk) + 1))
columns = list(schema) if columns is None else columns
meta = _meta_from_dtypes(columns, schema, [], [])
graph = HighLevelGraph({output_name: layer}, {output_name: set()})
return DataFrame(graph, output_name, meta, [None] * (len(parts) + 1))
12 changes: 7 additions & 5 deletions dask/dataframe/io/parquet/arrow.py
Original file line number Diff line number Diff line change
Expand Up @@ -1472,11 +1472,13 @@ def _process_metadata(
filters,
)

# Check if we need to pass a fragment for each output partition
read_from_paths = read_from_paths or False
# Check if we need to pass a fragment for each output partition.
# By default, we will avoid passing fragments in the graph unless
# the user has specified `read_from_paths=False`
partitions = partition_info.get("partitions", None)
pass_frags = (
filters
and (not read_from_paths)
and (read_from_paths is False)
and _need_fragments(filters, partition_info.get("partition_keys", None))
)

Expand All @@ -1492,7 +1494,7 @@ def _process_metadata(
make_part_kwargs={
"fs": fs,
"partition_keys": partition_info.get("partition_keys", None),
"partition_obj": partition_info.get("partitions", None),
"partition_obj": partitions,
"data_path": data_path,
"frag_map": frag_map if pass_frags else None,
},
Expand All @@ -1501,7 +1503,7 @@ def _process_metadata(
# Add common kwargs
common_kwargs = {
"partitioning": partition_info["partitioning"],
"partitions": partition_info["partitions"],
"partitions": partitions,
"categories": categories,
"filters": filters,
"schema": schema,
Expand Down
Loading