Update arrow to reduce plasma IPCs. (#3497)

This commit is contained in:
Philipp Moritz
2018-12-14 23:49:37 -05:00
committed by Robert Nishihara
parent fcc37021b2
commit b3bf608608
18 changed files with 27 additions and 34 deletions
+2 -2
View File
@@ -3,10 +3,10 @@
namespace ray {
ObjectBufferPool::ObjectBufferPool(const std::string &store_socket_name,
uint64_t chunk_size, int release_delay)
uint64_t chunk_size)
: default_chunk_size_(chunk_size) {
store_socket_name_ = store_socket_name;
ARROW_CHECK_OK(store_client_.Connect(store_socket_name_.c_str(), "", release_delay));
ARROW_CHECK_OK(store_client_.Connect(store_socket_name_.c_str(), ""));
}
ObjectBufferPool::~ObjectBufferPool() {
+1 -4
View File
@@ -40,10 +40,7 @@ class ObjectBufferPool {
/// \param store_socket_name The socket name of the store to which plasma clients
/// connect.
/// \param chunk_size The chunk size into which objects are to be split.
/// \param release_delay The number of release calls before objects are released
/// from the store client (FIFO).
ObjectBufferPool(const std::string &store_socket_name, const uint64_t chunk_size,
const int release_delay);
ObjectBufferPool(const std::string &store_socket_name, const uint64_t chunk_size);
~ObjectBufferPool();
+1 -4
View File
@@ -14,10 +14,7 @@ ObjectManager::ObjectManager(asio::io_service &main_service,
: config_(config),
object_directory_(std::move(object_directory)),
store_notification_(main_service, config_.store_socket_name),
// release_delay of 2 * config_.max_sends is to ensure the pool does not release
// an object prematurely whenever we reach the maximum number of sends.
buffer_pool_(config_.store_socket_name, config_.object_chunk_size,
/*release_delay=*/2 * config_.max_sends),
buffer_pool_(config_.store_socket_name, config_.object_chunk_size),
send_work_(send_service_),
receive_work_(receive_service_),
connection_pool_(),
@@ -15,8 +15,7 @@ namespace ray {
ObjectStoreNotificationManager::ObjectStoreNotificationManager(
boost::asio::io_service &io_service, const std::string &store_socket_name)
: store_client_(), socket_(io_service) {
ARROW_CHECK_OK(store_client_.Connect(store_socket_name.c_str(), "",
plasma::kPlasmaDefaultReleaseDelay));
ARROW_CHECK_OK(store_client_.Connect(store_socket_name.c_str(), ""));
ARROW_CHECK_OK(store_client_.Subscribe(&c_socket_));
boost::system::error_code ec;
@@ -154,8 +154,8 @@ class TestObjectManagerBase : public ::testing::Test {
server2.reset(new MockServer(main_service, om_config_2, gcs_client_2));
// connect to stores.
ARROW_CHECK_OK(client1.Connect(store_id_1, "", plasma::kPlasmaDefaultReleaseDelay));
ARROW_CHECK_OK(client2.Connect(store_id_2, "", plasma::kPlasmaDefaultReleaseDelay));
ARROW_CHECK_OK(client1.Connect(store_id_1, ""));
ARROW_CHECK_OK(client2.Connect(store_id_2, ""));
}
void TearDown() {
@@ -139,8 +139,8 @@ class TestObjectManagerBase : public ::testing::Test {
server2.reset(new MockServer(main_service, om_config_2, gcs_client_2));
// connect to stores.
ARROW_CHECK_OK(client1.Connect(store_id_1, "", plasma::kPlasmaDefaultReleaseDelay));
ARROW_CHECK_OK(client2.Connect(store_id_2, "", plasma::kPlasmaDefaultReleaseDelay));
ARROW_CHECK_OK(client1.Connect(store_id_1, ""));
ARROW_CHECK_OK(client2.Connect(store_id_2, ""));
}
void TearDown() {
@@ -74,8 +74,8 @@ class TestObjectManagerBase : public ::testing::Test {
GetNodeManagerConfig("raylet_2", store_sock_2), om_config_2, gcs_client_2));
// connect to stores.
ARROW_CHECK_OK(client1.Connect(store_sock_1, "", plasma::kPlasmaDefaultReleaseDelay));
ARROW_CHECK_OK(client2.Connect(store_sock_2, "", plasma::kPlasmaDefaultReleaseDelay));
ARROW_CHECK_OK(client1.Connect(store_sock_1, ""));
ARROW_CHECK_OK(client2.Connect(store_sock_2, ""));
}
void TearDown() {