mirror of
https://github.com/wassname/ray.git
synced 2026-08-12 12:20:11 +08:00
Update arrow and parquet-cpp. (#1875)
* Update arrow. * Fix bug. * Cherry-pick commit for fixing parquet segfault. * Update arrow and revert auto-releasing buffer commit. * Remove parquet cherry-pick.
This commit is contained in:
committed by
Philipp Moritz
parent
39cf6ff6e1
commit
d0fffec2d0
@@ -802,7 +802,7 @@ void process_transfer_request(event_loop *loop,
|
||||
/* We pass in 0 to indicate that the command should return immediately. */
|
||||
ARROW_CHECK_OK(
|
||||
conn->manager_state->plasma_conn->Get(&object_id, 1, 0, &object_buffer));
|
||||
if (object_buffer.data_size == -1) {
|
||||
if (object_buffer.data == nullptr) {
|
||||
/* If the object wasn't locally available, exit immediately. If the object
|
||||
* later appears locally, the requesting plasma manager should request the
|
||||
* transfer again. */
|
||||
@@ -823,15 +823,15 @@ void process_transfer_request(event_loop *loop,
|
||||
}
|
||||
|
||||
RAY_CHECK(object_buffer.metadata->data() ==
|
||||
object_buffer.data->data() + object_buffer.data_size);
|
||||
object_buffer.data->data() + object_buffer.data->size());
|
||||
PlasmaRequestBuffer *buf = new PlasmaRequestBuffer();
|
||||
buf->type = MessageType_PlasmaDataReply;
|
||||
buf->object_id = obj_id;
|
||||
/* We treat buf->data as a pointer to the concatenated data and metadata, so
|
||||
* we don't actually use buf->metadata. */
|
||||
buf->data = const_cast<uint8_t *>(object_buffer.data->data());
|
||||
buf->data_size = object_buffer.data_size;
|
||||
buf->metadata_size = object_buffer.metadata_size;
|
||||
buf->data_size = object_buffer.data->size();
|
||||
buf->metadata_size = object_buffer.metadata->size();
|
||||
|
||||
manager_conn->transfer_queue.push_back(buf);
|
||||
manager_conn->pending_object_transfers[object_id] = buf;
|
||||
|
||||
@@ -137,7 +137,7 @@ TEST plasma_nonblocking_get_tests(void) {
|
||||
|
||||
/* Test for object non-existence. */
|
||||
ARROW_CHECK_OK(client.Get(oid_array, 1, 0, &obj_buffer));
|
||||
ASSERT(obj_buffer.data_size == -1);
|
||||
ASSERT(obj_buffer.data == nullptr);
|
||||
|
||||
/* Test for the object being in local Plasma store. */
|
||||
/* First create object. */
|
||||
@@ -240,7 +240,7 @@ TEST plasma_get_tests(void) {
|
||||
PLASMA_DEFAULT_RELEASE_DELAY));
|
||||
ObjectID oid1 = ObjectID::from_random();
|
||||
ObjectID oid2 = ObjectID::from_random();
|
||||
ObjectBuffer obj_buffer;
|
||||
ObjectBuffer obj_buffer1;
|
||||
|
||||
ObjectID oid_array1[1] = {oid1};
|
||||
ObjectID oid_array2[1] = {oid2};
|
||||
@@ -254,17 +254,18 @@ TEST plasma_get_tests(void) {
|
||||
init_data_123(data->mutable_data(), data_size, 1);
|
||||
ARROW_CHECK_OK(client1.Seal(oid1));
|
||||
|
||||
ARROW_CHECK_OK(client1.Get(oid_array1, 1, -1, &obj_buffer));
|
||||
ASSERT(data->data()[0] == obj_buffer.data->data()[0]);
|
||||
ARROW_CHECK_OK(client1.Get(oid_array1, 1, -1, &obj_buffer1));
|
||||
ASSERT(data->data()[0] == obj_buffer1.data->data()[0]);
|
||||
|
||||
ObjectBuffer obj_buffer2;
|
||||
ARROW_CHECK_OK(
|
||||
client2.Create(oid2, data_size, metadata, metadata_size, &data));
|
||||
init_data_123(data->mutable_data(), data_size, 2);
|
||||
ARROW_CHECK_OK(client2.Seal(oid2));
|
||||
|
||||
ARROW_CHECK_OK(client1.Fetch(1, oid_array2));
|
||||
ARROW_CHECK_OK(client1.Get(oid_array2, 1, -1, &obj_buffer));
|
||||
ASSERT(data->data()[0] == obj_buffer.data->data()[0]);
|
||||
ARROW_CHECK_OK(client1.Get(oid_array2, 1, -1, &obj_buffer2));
|
||||
ASSERT(data->data()[0] == obj_buffer2.data->data()[0]);
|
||||
|
||||
sleep(1);
|
||||
ARROW_CHECK_OK(client1.Disconnect());
|
||||
|
||||
Reference in New Issue
Block a user