join#
- Dataset.join(ds: Dataset, join_type: str, num_partitions: int, on: Tuple[str] = ('id',), right_on: Tuple[str] | None = None, left_suffix: str | None = None, right_suffix: str | None = None, *, partition_size_hint: int | None = None, aggregator_ray_remote_args: Dict[str, Any] | None = None, validate_schemas: bool = False) Dataset[source]#
Join
Datasetson join keys.- Parameters:
ds – Other dataset to join against
join_type – The kind of join that should be performed, one of (“inner”, “left_outer”, “right_outer”, “full_outer”, “left_semi”, “right_semi”, “left_anti”, “right_anti”).
num_partitions – Total number of “partitions” input sequences will be split into with each partition being joined independently. Increasing number of partitions allows to reduce individual partition size, hence reducing memory requirements when individual partitions are being joined. Note that, consequently, this will also be a total number of blocks that will be produced as a result of executing join.
on – The columns from the left operand that will be used as keys for the join operation.
right_on – The columns from the right operand that will be used as keys for the join operation. When none,
onwill be assumed to be a list of columns to be used for the right dataset as well.left_suffix – (Optional) Suffix to be appended for columns of the left operand.
right_suffix – (Optional) Suffix to be appended for columns of the right operand.
partition_size_hint – (Optional) Deprecated and ignored. The join is now executed on the v2 hash-shuffle path, which sizes reduce-task memory from observed partition sizes rather than a hint. This parameter has no effect and will be removed in a future release.
aggregator_ray_remote_args – (Optional) Parameter overriding
ray.remoteargs passed when constructing joining (aggregator) workers.validate_schemas – (Optional) Controls whether validation of provided configuration against input schemas will be performed (defaults to false, since obtaining schemas could be prohibitively expensive).
- Returns:
A
Datasetthat holds rows of input left Dataset joined with the right Dataset based on join type and keys.
Note
This operation requires all inputs to be materialized in object store for it to execute.
Examples:
doubles_ds = ray.data.range(4).map( lambda row: {"id": row["id"], "double": int(row["id"]) * 2} ) squares_ds = ray.data.range(4).map( lambda row: {"id": row["id"], "square": int(row["id"]) ** 2} ) # Inner join example joined_ds = doubles_ds.join( squares_ds, join_type="inner", num_partitions=2, on=("id",), ) print(sorted(joined_ds.take_all(), key=lambda item: item["id"]))
[ {'id': 0, 'double': 0, 'square': 0}, {'id': 1, 'double': 2, 'square': 1}, {'id': 2, 'double': 4, 'square': 4}, {'id': 3, 'double': 6, 'square': 9} ]# Left anti-join example: find rows in doubles_ds that don't match squares_ds partial_squares_ds = ray.data.range(2).map( lambda row: {"id": row["id"] + 2, "square": int(row["id"]) ** 2} ) anti_joined_ds = doubles_ds.join( partial_squares_ds, join_type="left_anti", num_partitions=2, on=("id",), ) print(sorted(anti_joined_ds.take_all(), key=lambda item: item["id"]))
[ {'id': 0, 'double': 0}, {'id': 1, 'double': 2} ]# Left semi-join example: find rows in doubles_ds that have matches in squares_ds # (only returns columns from left dataset) semi_joined_ds = doubles_ds.join( squares_ds, join_type="left_semi", num_partitions=2, on=("id",), ) print(sorted(semi_joined_ds.take_all(), key=lambda item: item["id"]))
[ {'id': 0, 'double': 0}, {'id': 1, 'double': 2}, {'id': 2, 'double': 4}, {'id': 3, 'double': 6} ]