rgw: decoupling the bucket replication log from the bucket index
The bucket index logs used for multisite replication are currently stored in omap on the bucket index shards, along with the rest of the bucket index entries. Storing them in the index was a natural choice, because cls_rgw can write these log entries atomically when completing a bucket index transaction. To replicate a bucket, other zones process each of its bucket index shard logs independently, and store sync status markers with their position in each shard. This tight coupling between the replication strategy and the bucket's sharding scheme is the main challenge to supporting bucket resharding in multisite, because shuffling these log entries to a new set of shards would invalidate the sync status markers stored in other zones. My proposal, then, is to move the replication logs out of bucket index shards into a single log per bucket, and extend the consistency model to make up for the lack of atomic writes that we get from cls_rgw. The existing consistency model for object writes involves a) calling cls_rgw to prepare a bucket index transaction, b) writing the object's head to the data pool, then c) calling cls_rgw to complete the transaction. Since the write in b) is what makes the object visible to GET requests, we can reply to the client without waiting for c) to finish. If either b) or c) fails, the next bucket listing will find an entry that was prepared but not completed, and we'll check whether the head object exists and use the 'dir suggest' call to update the bucket index accordingly. If we move the replication log to a separate object, we'll need to write to that as well before completing the transaction. And when dir suggest finds head objects for uncompleted transactions, it can (re)write their replication log entries before updating the bucket index. This recovery means that we can still reply to the client before writing to the replication log, so the client won't see any extra latency. This change also gives us the opportunity to move away from omap and the challenges associated with trimming. Yehuda wrote cls_fifo in https://github.com/ceph/ceph/pull/30797 with the datalog in mind, and that could be a good fit for these bucket replication logs as well.
On Mon, 9 Dec 2019, Casey Bodley wrote:
The bucket index logs used for multisite replication are currently stored in omap on the bucket index shards, along with the rest of the bucket index entries. Storing them in the index was a natural choice, because cls_rgw can write these log entries atomically when completing a bucket index transaction.
To replicate a bucket, other zones process each of its bucket index shard logs independently, and store sync status markers with their position in each shard. This tight coupling between the replication strategy and the bucket's sharding scheme is the main challenge to supporting bucket resharding in multisite, because shuffling these log entries to a new set of shards would invalidate the sync status markers stored in other zones.
My proposal, then, is to move the replication logs out of bucket index shards into a single log per bucket, and extend the consistency model to make up for the lack of atomic writes that we get from cls_rgw.
The existing consistency model for object writes involves a) calling cls_rgw to prepare a bucket index transaction, b) writing the object's head to the data pool, then c) calling cls_rgw to complete the transaction. Since the write in b) is what makes the object visible to GET requests, we can reply to the client without waiting for c) to finish. If either b) or c) fails, the next bucket listing will find an entry that was prepared but not completed, and we'll check whether the head object exists and use the 'dir suggest' call to update the bucket index accordingly.
If we move the replication log to a separate object, we'll need to write to that as well before completing the transaction. And when dir suggest finds head objects for uncompleted transactions, it can (re)write their replication log entries before updating the bucket index. This recovery means that we can still reply to the client before writing to the replication log, so the client won't see any extra latency.
This makes sense to me! Just to make sure I understand: (a) prepare the bucket index txn (b) update the head (c) write the replication log entry (d) clean up the index txn This means that if we fail after b and the dir_suggest replays, then we may get duplicated (c) items. Does it also mean that we might not notice the dropped replication log entry right away? Or maybe the multisite map that tells us which buckets may be dirty means we can check those bucket indexes for any possible in-progress transaction? Otherwise we might end up not registring the replication log item until (much) later. This also means that there could be duplicate items in the replication log for the same update. An alternative might be to do steps (a) and (c) in parallel, but then the replication log entry might reflect a head update that hasn't updated yet (or perhaps never happens), which would make the replication machinery more complex.
This change also gives us the opportunity to move away from omap and the challenges associated with trimming. Yehuda wrote cls_fifo in https://github.com/ceph/ceph/pull/30797 with the datalog in mind, and that could be a good fit for these bucket replication logs as well.
+1 sage
On Mon, Dec 9, 2019 at 11:44 PM Sage Weil <sage@newdream.net> wrote:
On Mon, 9 Dec 2019, Casey Bodley wrote:
The bucket index logs used for multisite replication are currently stored in omap on the bucket index shards, along with the rest of the bucket index entries. Storing them in the index was a natural choice, because cls_rgw can write these log entries atomically when completing a bucket index transaction.
To replicate a bucket, other zones process each of its bucket index shard logs independently, and store sync status markers with their position in each shard. This tight coupling between the replication strategy and the bucket's sharding scheme is the main challenge to supporting bucket resharding in multisite, because shuffling these log entries to a new set of shards would invalidate the sync status markers stored in other zones.
My proposal, then, is to move the replication logs out of bucket index shards into a single log per bucket, and extend the consistency model to make up for the lack of atomic writes that we get from cls_rgw.
The existing consistency model for object writes involves a) calling cls_rgw to prepare a bucket index transaction, b) writing the object's head to the data pool, then c) calling cls_rgw to complete the transaction. Since the write in b) is what makes the object visible to GET requests, we can reply to the client without waiting for c) to finish. If either b) or c) fails, the next bucket listing will find an entry that was prepared but not completed, and we'll check whether the head object exists and use the 'dir suggest' call to update the bucket index accordingly.
If we move the replication log to a separate object, we'll need to write to that as well before completing the transaction. And when dir suggest finds head objects for uncompleted transactions, it can (re)write their replication log entries before updating the bucket index. This recovery means that we can still reply to the client before writing to the replication log, so the client won't see any extra latency.
This makes sense to me! Just to make sure I understand:
(a) prepare the bucket index txn (b) update the head (c) write the replication log entry (d) clean up the index txn
This means that if we fail after b and the dir_suggest replays, then we may get duplicated (c) items. Does it also mean that we might not notice the dropped replication log entry right away? Or maybe the multisite map that tells us which buckets may be dirty means we can check those bucket indexes for any possible in-progress transaction? Otherwise we might end up not registring the replication log item until (much) later.
I wouldn't be worried about duplicate items, in the worst case we'd try to sync the same entry twice, but we'd identify it as already existing (same would work for deletes). There is no efficient way currently that would let us check for in-flight transactions at the bucket index.
This also means that there could be duplicate items in the replication log for the same update.
An alternative might be to do steps (a) and (c) in parallel, but then the replication log entry might reflect a head update that hasn't updated yet (or perhaps never happens), which would make the replication machinery more complex.
I'm not sure how that could be solved, would be inherently racy. Specifically there'd be issue with deletes that would be hard to solve without introducing tombstones. Yehuda
This change also gives us the opportunity to move away from omap and the challenges associated with trimming. Yehuda wrote cls_fifo in https://github.com/ceph/ceph/pull/30797 with the datalog in mind, and that could be a good fit for these bucket replication logs as well.
+1
sage _______________________________________________ Dev mailing list -- dev@ceph.io To unsubscribe send an email to dev-leave@ceph.io
On 12/9/19 4:44 PM, Sage Weil wrote:
On Mon, 9 Dec 2019, Casey Bodley wrote:
The bucket index logs used for multisite replication are currently stored in omap on the bucket index shards, along with the rest of the bucket index entries. Storing them in the index was a natural choice, because cls_rgw can write these log entries atomically when completing a bucket index transaction.
To replicate a bucket, other zones process each of its bucket index shard logs independently, and store sync status markers with their position in each shard. This tight coupling between the replication strategy and the bucket's sharding scheme is the main challenge to supporting bucket resharding in multisite, because shuffling these log entries to a new set of shards would invalidate the sync status markers stored in other zones.
My proposal, then, is to move the replication logs out of bucket index shards into a single log per bucket, and extend the consistency model to make up for the lack of atomic writes that we get from cls_rgw.
The existing consistency model for object writes involves a) calling cls_rgw to prepare a bucket index transaction, b) writing the object's head to the data pool, then c) calling cls_rgw to complete the transaction. Since the write in b) is what makes the object visible to GET requests, we can reply to the client without waiting for c) to finish. If either b) or c) fails, the next bucket listing will find an entry that was prepared but not completed, and we'll check whether the head object exists and use the 'dir suggest' call to update the bucket index accordingly.
If we move the replication log to a separate object, we'll need to write to that as well before completing the transaction. And when dir suggest finds head objects for uncompleted transactions, it can (re)write their replication log entries before updating the bucket index. This recovery means that we can still reply to the client before writing to the replication log, so the client won't see any extra latency. This makes sense to me! Just to make sure I understand:
(a) prepare the bucket index txn (b) update the head (c) write the replication log entry (d) clean up the index txn
This means that if we fail after b and the dir_suggest replays, then we may get duplicated (c) items. Does it also mean that we might not notice the dropped replication log entry right away? Or maybe the multisite map that tells us which buckets may be dirty means we can check those bucket indexes for any possible in-progress transaction? Otherwise we might end up not registring the replication log item until (much) later.
That's true, dir suggest only guarantees that the next bucket listing is consistent; if nobody lists the bucket, then this recovery never runs. We're in the same situation now, where cls_rgw's rgw_dir_suggest_changes() call is writing to the bilog.
This also means that there could be duplicate items in the replication log for the same update.
An alternative might be to do steps (a) and (c) in parallel, but then the replication log entry might reflect a head update that hasn't updated yet (or perhaps never happens), which would make the replication machinery more complex.
This change also gives us the opportunity to move away from omap and the challenges associated with trimming. Yehuda wrote cls_fifo in https://github.com/ceph/ceph/pull/30797 with the datalog in mind, and that could be a good fit for these bucket replication logs as well. +1
sage _______________________________________________ Dev mailing list -- dev@ceph.io To unsubscribe send an email to dev-leave@ceph.io
On Tue, Dec 10, 2019 at 3:05 AM Casey Bodley <cbodley@redhat.com> wrote:
The bucket index logs used for multisite replication are currently stored in omap on the bucket index shards, along with the rest of the bucket index entries. Storing them in the index was a natural choice, because cls_rgw can write these log entries atomically when completing a bucket index transaction.
To replicate a bucket, other zones process each of its bucket index shard logs independently, and store sync status markers with their position in each shard. This tight coupling between the replication strategy and the bucket's sharding scheme is the main challenge to supporting bucket resharding in multisite, because shuffling these log entries to a new set of shards would invalidate the sync status markers stored in other zones.
Can we alternatively make the sync status markers resilient to changes introduced by bucket resharding change? I don't have a specific method in mind, but presuming that it should be relatively simple, risk-free and backward compatible compared to making changes in the write path. Thanks, K.Prasad -- *-----------------------------------------------------------------------------------------* *This email and any files transmitted with it are confidential and intended solely for the use of the individual or entity to whom they are addressed. If you have received this email in error, please notify the system manager. This message contains confidential information and is intended only for the individual named. If you are not the named addressee, you should not disseminate, distribute or copy this email. Please notify the sender immediately by email if you have received this email by mistake and delete this email from your system. If you are not the intended recipient, you are notified that disclosing, copying, distributing or taking any action in reliance on the contents of this information is strictly prohibited.***** **** *Any views or opinions presented in this email are solely those of the author and do not necessarily represent those of the organization. Any information on shares, debentures or similar instruments, recommended product pricing, valuations and the like are for information purposes only. It is not meant to be an instruction or recommendation, as the case may be, to buy or to sell securities, products, services nor an offer to buy or sell securities, products or services unless specifically stated to be so on behalf of the Flipkart group. Employees of the Flipkart group of companies are expressly required not to make defamatory statements and not to infringe or authorise any infringement of copyright or any other legal right by email communications. Any such communication is contrary to organizational policy and outside the scope of the employment of the individual concerned. The organization will not accept any liability in respect of such communication, and the employee responsible will be personally liable for any damages or other liability arising.***** **** *Our organization accepts no liability for the content of this email, or for the consequences of any actions taken on the basis of the information *provided,* unless that information is subsequently confirmed in writing. If you are not the intended recipient, you are notified that disclosing, copying, distributing or taking any action in reliance on the contents of this information is strictly prohibited.* _-----------------------------------------------------------------------------------------_
On 12/9/19 10:03 PM, Prasad Krishnan wrote:
On Tue, Dec 10, 2019 at 3:05 AM Casey Bodley <cbodley@redhat.com <mailto:cbodley@redhat.com>> wrote:
The bucket index logs used for multisite replication are currently stored in omap on the bucket index shards, along with the rest of the bucket index entries. Storing them in the index was a natural choice, because cls_rgw can write these log entries atomically when completing a bucket index transaction.
To replicate a bucket, other zones process each of its bucket index shard logs independently, and store sync status markers with their position in each shard. This tight coupling between the replication strategy and the bucket's sharding scheme is the main challenge to supporting bucket resharding in multisite, because shuffling these log entries to a new set of shards would invalidate the sync status markers stored in other zones.
Can we alternatively make the sync status markers resilient to changes introduced by bucket resharding change? I don't have a specific method in mind, but presuming that it should be relatively simple, risk-free and backward compatible compared to making changes in the write path.
The omap keys of these bilog entries are essentially just encoding a version number that's local to that index shard (see get_index_ver_key() in cls_rgw.cc), so there's no defined ordering relationship between keys in different shards. Other zones store an array of sync status markers (ie omap keys) with their position in each log shard. If the bucket gets resharded, the size of this sync status array no longer matches the new bilog shards, with no way to associate its markers with the new sharding scheme.
Thanks, K.Prasad
/-----------------------------------------------------------------------------------------/
/This email and any files transmitted with it are confidential and intended solely for the use of the individual or entity to whom they are addressed. If you have received this email in error, please notify the system manager. This message contains confidential information and is intended only for the individual named. If you are not the named addressee, you should not disseminate, distribute or copy this email. Please notify the sender immediately by email if you have received this email by mistake and delete this email from your system. If you are not the intended recipient, you are notified that disclosing, copying, distributing or taking any action in reliance on the contents of this information is strictly prohibited./
/Any views or opinions presented in this email are solely those of the author and do not necessarily represent those of the organization. Any information on shares, debentures or similar instruments, recommended product pricing, valuations and the like are for information purposes only. It is not meant to be an instruction or recommendation, as the case may be, to buy or to sell securities, products, services nor an offer to buy or sell securities, products or services unless specifically stated to be so on behalf of the Flipkart group. Employees of the Flipkart group of companies are expressly required not to make defamatory statements and not to infringe or authorise any infringement of copyright or any other legal right by email communications. Any such communication is contrary to organizational policy and outside the scope of the employment of the individual concerned. The organization will not accept any liability in respect of such communication, and the employee responsible will be personally liable for any damages or other liability arising./
/Our organization accepts no liability for the content of this email, or for the consequences of any actions taken on the basis of the information /provided,/ unless that information is subsequently confirmed in writing. If you are not the intended recipient, you are notified that disclosing, copying, distributing or taking any action in reliance on the contents of this information is strictly prohibited./
/-----------------------------------------------------------------------------------------/
On Mon, Dec 9, 2019 at 11:35 PM Casey Bodley <cbodley@redhat.com> wrote:
The bucket index logs used for multisite replication are currently stored in omap on the bucket index shards, along with the rest of the bucket index entries. Storing them in the index was a natural choice, because cls_rgw can write these log entries atomically when completing a bucket index transaction.
To replicate a bucket, other zones process each of its bucket index shard logs independently, and store sync status markers with their position in each shard. This tight coupling between the replication strategy and the bucket's sharding scheme is the main challenge to supporting bucket resharding in multisite, because shuffling these log entries to a new set of shards would invalidate the sync status markers stored in other zones.
Note that the new bucket granularity work tackle part of the problem and lays foundation for solving it by managing unbalanced replication, where source bucket and destination bucket have different number of shards. So when a (target) bucket is resharded, we could still craft new markers for it to track its original source, even if they don't have the same number of shards (this is not implemented, but should be relatively easy to do). In the other direction, when source is resharded, currently the new bucket instance is handled as a new bucket at the target, so it triggers a full sync, but only actual new entries are being fetched. We can probably find a better and more optimal solution where we finish syncing the old entries from the old instance, and have the new one set to incremental sync from the start. That been said, I'm not against decoupling the logs and the index.
My proposal, then, is to move the replication logs out of bucket index shards into a single log per bucket, and extend the consistency model to make up for the lack of atomic writes that we get from cls_rgw.
Sharded log? If it's not sharded then you're going to introduce an object where all IO for that bucket will serialize over with high enough pressure (as I assume writes to it will be async).
The existing consistency model for object writes involves a) calling cls_rgw to prepare a bucket index transaction, b) writing the object's head to the data pool, then c) calling cls_rgw to complete the transaction. Since the write in b) is what makes the object visible to GET requests, we can reply to the client without waiting for c) to finish. If either b) or c) fails, the next bucket listing will find an entry that was prepared but not completed, and we'll check whether the head object exists and use the 'dir suggest' call to update the bucket index accordingly.
If we move the replication log to a separate object, we'll need to write to that as well before completing the transaction. And when dir suggest finds head objects for uncompleted transactions, it can (re)write their replication log entries before updating the bucket index. This recovery means that we can still reply to the client before writing to the replication log, so the client won't see any extra latency.
We'd need to write to that log after head has been written, otherwise clients reading it will race with the head creation. In which case I'm not sure how we don't introduce extra latency, since we can complete the transaction only after it has been written (or we'd potentially lose this log entry). There is the option to only write to the log (async) and then have a lazy process in a separate thread that goes over that log and completes the transactions, or hand it over to a separate workqueue.
This change also gives us the opportunity to move away from omap and the challenges associated with trimming. Yehuda wrote cls_fifo in https://github.com/ceph/ceph/pull/30797 with the datalog in mind, and that could be a good fit for these bucket replication logs as well.
Yeah, cls_fifo would be a good fit for it. Yehuda
_______________________________________________ Dev mailing list -- dev@ceph.io To unsubscribe send an email to dev-leave@ceph.io
Thanks for the feedback! On 12/10/19 1:52 AM, Yehuda Sadeh-Weinraub wrote:
The bucket index logs used for multisite replication are currently stored in omap on the bucket index shards, along with the rest of the bucket index entries. Storing them in the index was a natural choice, because cls_rgw can write these log entries atomically when completing a bucket index transaction.
To replicate a bucket, other zones process each of its bucket index shard logs independently, and store sync status markers with their position in each shard. This tight coupling between the replication strategy and the bucket's sharding scheme is the main challenge to supporting bucket resharding in multisite, because shuffling these log entries to a new set of shards would invalidate the sync status markers stored in other zones. Note that the new bucket granularity work tackle part of the problem and lays foundation for solving it by managing unbalanced replication, where source bucket and destination bucket have different number of shards. So when a (target) bucket is resharded, we could still craft new markers for it to track its original source, even if they don't have
On Mon, Dec 9, 2019 at 11:35 PM Casey Bodley <cbodley@redhat.com> wrote: the same number of shards (this is not implemented, but should be relatively easy to do). In the other direction, when source is resharded, currently the new bucket instance is handled as a new bucket at the target, so it triggers a full sync, but only actual new entries are being fetched. We can probably find a better and more optimal solution where we finish syncing the old entries from the old instance, and have the new one set to incremental sync from the start.
Yeah, I'd like to find an approach that avoids full sync on reshard, because those buckets will tend to be the big ones. But I also don't want to end up in the same situation as bucket deletion, where we can't delete the old bucket index shards until all other zones finish processing its logs.
That been said, I'm not against decoupling the logs and the index.
My proposal, then, is to move the replication logs out of bucket index shards into a single log per bucket, and extend the consistency model to make up for the lack of atomic writes that we get from cls_rgw. Sharded log? If it's not sharded then you're going to introduce an object where all IO for that bucket will serialize over with high enough pressure (as I assume writes to it will be async).
I agree that the log write latency is important for scalability. It won't impact PutObj performance directly because it's async, but it will increase the average time delta between the index prepares and completes, so bucket listings will tend to do more recovery work with dir suggest. Sharding wouldn't be my first choice though. It can reduce the write contention by a factor of num-shards, but we'd have to reintroduce dynamic sharding to scale this up for large buckets. Instead, we can batch up log writes in radosgw until the previous batch finishes. That limits write contention to the number of gateways, rather than the number of PutObj ops. I think batching will be important to get good performance out of cls_fifo anyway, where there's some overhead to discover the approximate append position, resend appends that land on full rados objects, etc. If we're batching up the log writes, it could also make sense to add a batch interface to cls_rgw for the index completions.
The existing consistency model for object writes involves a) calling cls_rgw to prepare a bucket index transaction, b) writing the object's head to the data pool, then c) calling cls_rgw to complete the transaction. Since the write in b) is what makes the object visible to GET requests, we can reply to the client without waiting for c) to finish. If either b) or c) fails, the next bucket listing will find an entry that was prepared but not completed, and we'll check whether the head object exists and use the 'dir suggest' call to update the bucket index accordingly.
If we move the replication log to a separate object, we'll need to write to that as well before completing the transaction. And when dir suggest finds head objects for uncompleted transactions, it can (re)write their replication log entries before updating the bucket index. This recovery means that we can still reply to the client before writing to the replication log, so the client won't see any extra latency. We'd need to write to that log after head has been written, otherwise clients reading it will race with the head creation. In which case I'm not sure how we don't introduce extra latency, since we can complete the transaction only after it has been written (or we'd potentially lose this log entry). There is the option to only write to the log (async) and then have a lazy process in a separate thread that goes over that log and completes the transactions, or hand it over to a separate workqueue.
These can still be asynchronous with respect to the http response, as long as the log writes succeed before we complete the index transactions. For example, by having the log write's AioCompletion callback schedule the index completion op. Batching makes this more complicated, but still has to provide the same guarantee.
This change also gives us the opportunity to move away from omap and the challenges associated with trimming. Yehuda wrote cls_fifo in https://github.com/ceph/ceph/pull/30797 with the datalog in mind, and that could be a good fit for these bucket replication logs as well. Yeah, cls_fifo would be a good fit for it.
Yehuda
_______________________________________________ Dev mailing list -- dev@ceph.io To unsubscribe send an email to dev-leave@ceph.io
participants (4)
-
Casey Bodley
-
Prasad Krishnan
-
Sage Weil
-
Yehuda Sadeh-Weinraub