mirror of
https://github.com/wassname/ray.git
synced 2026-07-23 13:10:11 +08:00
Change Python's ObjectID to ObjectRef (#9353)
This commit is contained in:
@@ -119,7 +119,7 @@ def test_global_state_api(shutdown_only):
|
||||
|
||||
# A driver/worker creates a temporary object during startup. Although the
|
||||
# temporary object is freed immediately, in a rare case, we can still find
|
||||
# the object ID in GCS because Raylet removes the object ID from GCS
|
||||
# the object ref in GCS because Raylet removes the object ref from GCS
|
||||
# asynchronously.
|
||||
# Because we can't control when workers create the temporary objects, so
|
||||
# We can't assert that `ray.objects()` returns an empty dict. Here we just
|
||||
@@ -281,22 +281,22 @@ def test_specific_job_id():
|
||||
ray.shutdown()
|
||||
|
||||
|
||||
def test_object_id_properties():
|
||||
def test_object_ref_properties():
|
||||
id_bytes = b"00112233445566778899"
|
||||
object_id = ray.ObjectID(id_bytes)
|
||||
assert object_id.binary() == id_bytes
|
||||
object_id = ray.ObjectID.nil()
|
||||
assert object_id.is_nil()
|
||||
object_ref = ray.ObjectRef(id_bytes)
|
||||
assert object_ref.binary() == id_bytes
|
||||
object_ref = ray.ObjectRef.nil()
|
||||
assert object_ref.is_nil()
|
||||
with pytest.raises(ValueError, match=r".*needs to have length 20.*"):
|
||||
ray.ObjectID(id_bytes + b"1234")
|
||||
ray.ObjectRef(id_bytes + b"1234")
|
||||
with pytest.raises(ValueError, match=r".*needs to have length 20.*"):
|
||||
ray.ObjectID(b"0123456789")
|
||||
object_id = ray.ObjectID.from_random()
|
||||
assert not object_id.is_nil()
|
||||
assert object_id.binary() != id_bytes
|
||||
id_dumps = pickle.dumps(object_id)
|
||||
ray.ObjectRef(b"0123456789")
|
||||
object_ref = ray.ObjectRef.from_random()
|
||||
assert not object_ref.is_nil()
|
||||
assert object_ref.binary() != id_bytes
|
||||
id_dumps = pickle.dumps(object_ref)
|
||||
id_from_dumps = pickle.loads(id_dumps)
|
||||
assert id_from_dumps == object_id
|
||||
assert id_from_dumps == object_ref
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
@@ -434,7 +434,7 @@ def test_pandas_parquet_serialization():
|
||||
|
||||
|
||||
def test_socket_dir_not_existing(shutdown_only):
|
||||
random_name = ray.ObjectID.from_random().hex()
|
||||
random_name = ray.ObjectRef.from_random().hex()
|
||||
temp_raylet_socket_dir = os.path.join(ray.utils.get_ray_temp_dir(),
|
||||
"tests", random_name)
|
||||
temp_raylet_socket_name = os.path.join(temp_raylet_socket_dir,
|
||||
@@ -487,20 +487,20 @@ def test_put_pins_object(ray_start_object_store_memory):
|
||||
obj = np.ones(200 * 1024, dtype=np.uint8)
|
||||
x_id = ray.put(obj)
|
||||
x_binary = x_id.binary()
|
||||
assert (ray.get(ray.ObjectID(x_binary)) == obj).all()
|
||||
assert (ray.get(ray.ObjectRef(x_binary)) == obj).all()
|
||||
|
||||
# x cannot be evicted since x_id pins it
|
||||
for _ in range(10):
|
||||
ray.put(np.zeros(10 * 1024 * 1024))
|
||||
assert (ray.get(x_id) == obj).all()
|
||||
assert (ray.get(ray.ObjectID(x_binary)) == obj).all()
|
||||
assert (ray.get(ray.ObjectRef(x_binary)) == obj).all()
|
||||
|
||||
# now it can be evicted since x_id pins it but x_binary does not
|
||||
del x_id
|
||||
for _ in range(10):
|
||||
ray.put(np.zeros(10 * 1024 * 1024))
|
||||
assert not ray.worker.global_worker.core_worker.object_exists(
|
||||
ray.ObjectID(x_binary))
|
||||
ray.ObjectRef(x_binary))
|
||||
|
||||
# weakref put
|
||||
y_id = ray.put(obj, weakref=True)
|
||||
@@ -530,7 +530,7 @@ def test_decorated_function(ray_start_regular):
|
||||
|
||||
|
||||
def test_get_postprocess(ray_start_regular):
|
||||
def get_postprocessor(object_ids, values):
|
||||
def get_postprocessor(object_refs, values):
|
||||
return [value for value in values if value > 0]
|
||||
|
||||
ray.worker.global_worker._post_get_hooks.append(get_postprocessor)
|
||||
@@ -657,10 +657,10 @@ def test_lease_request_leak(shutdown_only):
|
||||
# from the raylet.
|
||||
tasks = []
|
||||
for _ in range(10):
|
||||
oid = ray.put(1)
|
||||
obj_ref = ray.put(1)
|
||||
for _ in range(2):
|
||||
tasks.append(f.remote(oid))
|
||||
del oid
|
||||
tasks.append(f.remote(obj_ref))
|
||||
del obj_ref
|
||||
ray.get(tasks)
|
||||
|
||||
time.sleep(
|
||||
|
||||
Reference in New Issue
Block a user