9 ms·
I just finished rolling out Pulsar to 8 AWS regions with geo-replication. Messages rates are currently at about 50k msgs/sec but still in the process of migrati
by addisonj 7y ago
I just finished rolling out Pulsar to 8 AWS regions with geo-replication. Messages rates are currently at about 50k msgs/sec but still in the process of migrating many more applications. We run on top of kubernetes (EKS).
It took about 5 months for our implementation with a chunk of that work mostly about figuring out how to integrate our internal auth as well as a using hashicorp vault as a clean automated way to get auth tokens for an AWS IAM role.
Overall, we are very pleased and the rest of the engineering org is very excited about it and planning to migrate most of our SQS and Kinesis apps.
Ask me anything in thread and will try and answer questions. At some point we will do a blog post on our experience.
- rocky1138 7y agoWhy did you choose Pulsar?
- addisonj 7y agoThe main driver for Pulsar is that we have a number of different messaging use cases, some more "pub/sub" like and some that are more "log" like. Pulsar really does unify those two worlds while also being a ton more flexible than any hosted options. For example, Kinesis is really limiting with the limited retention and making it very difficult to do any real ordering at scale due to the really tiny size of each shard. Similarly, SQS does pub/sub well, but we keep finding that we do need to use the data more than the first initial delivery. Instead of having multiple systems where we store that data we have one. As for why we didn't go with Kafka, the biggest single reason is that Pulsar is easier operationally with no needing to re-balance and also with the awesome feature that is tiered storage via offloading that allows us to actually do topics that have unlimited retention. Perhaps more importantly for the adoption though is pub/sub is much easier with Pulsar and the API is just much easier to reason about for developers than all the complexity of consumer groups, etc. There are a ton of other nice things like being able to have topics be so cheap such that we can have hundred of thousands and all of the built-in multi-tenancy features, geo-replication, flexible ACL system, pulsar functions and pulsar IO and many other things that really have us excited about all the capabilities
- dominotw 7y ago> able to have topics be so cheap For GDPR a lot of us has to do exportable 'user activity'. Can you in theory have a topic/user ( we had like 50 million users) and publish any user activity to that topic?
- addisonj 7y agoPulsar docs indicate "millions" of topics but IDK what 50 million would look like but from what I know I would be a bit nervous about it :) It might be worth chatting with Pulsar devs on their slack community (https://apache-pulsar.herokuapp.com/ https://apache-pulsar.herokuapp.com/). Most commonly what I hear people doing for this is either one of two approaches (or a combination of both): - encrypt the user data and delete the key, eventually the user data will get removed - regularly compact the topic (pulsar has a compaction feature) and write in a tombstone record which will remove any user data after compaction
- bubbleRefuge 7y agoIsn't using Kubernetes kind of an anti-pattern due to failover and rebalancing logic clashing? If Kubernetes is killing and re-starting nodes and the cluster's brokers are detecting dead brokers and rebalancing partitions as a result, it seems counterproductive.
- rhizome 7y agoTo the degree that that conflict exists in their implementation I would think that it's possible to account for all of that.
- addisonj 7y agoThis is one of the main benefits of Pulsar is that because state is split between brokers and bookkeeper and bookkeeper doesn't need re-balanced (due to it's segment based architecture where you choose new bookies with each new segment), we really don't have to worry about re-balancing (in general, not just in case of failover) of storage. It is true that topics map to a single broker, but generally, Pulsar has really good limits on memory so we don't see nodes getting killed by limits and we only really see re-scheduling for real issues. While there certainly is some aspects you need to be aware of, generally, Pulsar is much more "cloud native" and maps quiet well to k8s primitives.
- skube 7y agoUsing kubernetes is always an anti-pattern.
- jupp0r 7y agoIt's a good pattern because it regularly forces you to deal with pods/nodes going away so that the system is designed to handle this well without human intervention. There is no system where nodes don't go away because of hardware errors/updates/decommissions, so you might as well establish it as unexciting routine from the start.
- ypcx 7y agoAlrighty, a few questions: - what k8s definitions do you use, e.g. do you use the official Helm Chart, or have you written your .yaml's from scratch? - have you practiced disaster recovery scenarios in the context of k8s? Can you describe them briefly? - how do you upgrade/redeploy the Pulsar k8s components, i.e. does this cause the Bookies to trigger a cluster rebalance, or does it trigger the Autorecovery - for the Bookies, do you use AWS EBS volumes with the EKS or just local instance storage (that is, if you use persistent topics) - do you use the Proxy pod's EKS k8s pod IPs as exposed on the AWS network, or do you use a NodePort type of service for the Proxy components (using the EKS node IPs) - have you been bitten by the recent EKS k8s network plugin bug (loss of pod connectivity), and/or how do you maintain your EKS cluster - do you run your EKS nodes in a multi-AZ setting?
- addisonj 7y agofor the k8s definitions, we started with the helm chart, rendered the template, and then moved it into kustomize, as that is our tool of choice ATM, IDK if I would recommend that approach for everyone (we expect we might move to helm v3 at some point) but it was a good choice for us. We have practiced some disaster recovery, but it isn't 100% exhaustive (is it ever?), however it is also aided by how Pulsar is designed. We have killed bookie nodes as well as lost all our state in zookeeper. The first is pretty easily handled by the replication factor of bookkeeper data and for zookeeper we do extra backup step and just dump the state to s3 and can restore it. What we haven't tested in practice but now how to do theoretically is to restore a k8s stateful set from EBS volume snapshots. However, we see that as a real edge case. In Pulsar, we offload our data to s3 after a few hours, so we only need to worry about potentially losing a few hours of data in BK, as the zookeeper state is very easy to just snapshot and restore from s3. In other words, we are still working on getting more and more confident with data and don't yet recommend teams use it for mission critical non-recoverable data, but there are a ton of uses cases for it now and we can continue to improve on the DR front We have done multiple upgrades and deploy all the time. Because bookkeeper nodes are in a stateful set and we have don't do automated rollouts, we manually have a process to replace the BK nodes. However, they don't trigger a re-balance as it closes gracefully and then re-attaches the EBS volume from the stateful set We use EBS volumes, we use a piops volumes for the journal and a larger slower volume for the ledger store. THis is one of the great parts of bookkeeper design is that the two disks pools are separate so we just need a small chunk of really fast storage and then the journaled data is copied over to the ledger volume by a background process. We figure for really high write throughput we could use instance storage for the journal volume and EBS for ledger, but that would have some complications on recovery but still easier than having to rebuild the whole ledger data. We use the pulsar proxy and expose it via a k8s service with the AWS specific NLB annotations. We haven't had any issues with the k8s plugin and haven't really had any issues with EKS version upgrades. We just add new nodes when we migrate the kubelets Yes, we have automation (via terraform) to allow us to add many different pools of compute and we use labels and taints to get specific apps mapped to specific pools of compute. For Pulsar, we run all the components multi-AZ
- GordonS 7y agoReally interested why you chose Pulsar over RabbitMQ and others?
- deleted 7y ago[deleted]
- addisonj 7y agoI have used (and deployed) rabbitmq in the past and really love it for pub/sub, but for our needs, we keep needing retention, particularly long retention that we process with Flink for computing views. Having one system to do both is great for us.
- GordonS 7y agoSorry, I'd missed that Pulsar was a streaming log system (like kafka), as well as a pub/sub system. HN title misled me :)
- bubbleRefuge 7y agoHow does python stream processing work. Are the modules running in the JVM ? Jython ?
- matteomerli 7y agoThe user code it's all running in a native CPython interpreter.
- willvarfar 7y agoWhat's your plan on disaster recovery? Do your workers track their own cursors, and if so, how does that work across regions?
- addisonj 7y agoIn Pulsar, offsets are tracked by the service as part of the bookkeeper data (unless you use the reader API which is only really needed for advanced use cases like Flink), that means we just need to do DR for bookkeepers, which I touch on in another response but the tl:dr; is that we have a 3x replication factor as well as EBS snapshots
- 3fe9a03ccd14ca5 7y agoThe first question I have is why? SQS seems like such a simple thing to keep hosted.
- jjeaff 7y agoThey said they are currently doing 50k messages a second and they aren't even done migrating everything over. 50k messages a second would cost you around $50k a month for AWS sqs, (math could be wrong, didn't double check). Plus, with sqs, you get what they have. No customizations.
- rumanator 7y agoEven if you were off by an order of magnitude, that expenditure level is not justified to run a message broker service.
- jhh 7y ago5k a month would be an amazing deal.
- deanCommie 7y agoI sincerely doubt they are sustaining 50k msgs/second. Likely that's the MAXIMUM throughput. No way they would actually hit that sustained throughput for the entire month. Even the other justifications about wanting to reference messages after delivery do not to me justify migrating off SQS/Kinesis, especially not at cost of 5 months development effort.
- nostrebored 7y agoWhy is this hard to believe? Maybe they have chatty IoT devices. Maybe a ton of sensors monitoring manufacturing plants for multiple different metrics in real-time. Maybe they just have large scale.
- manigandham 7y agoWe do 100k msgs/sec minimum in our streaming platform. This is on the low end for adtech and other industries.
- mavdi 7y agoNot questioning your judgement but interested to know about the factors moving you away from Kinesis.
- ckdarby 7y agoKinesis is very expensive in the long run. There's almost always an intersection point on AWS where you need to consider moving away from AWS services/"managed services" and bring it in house.
- addisonj 7y agoBiggest pain points with Kinesis: - ordering is really hard, you don't get guaranted ordering unless you write one message at a time or do a lot of complexity on writes (see https://brandur.org/kinesis-order https://brandur.org/kinesis-order) and the shards are simply too small for many of our ordered use cases - cost, we just don't send some data right now because it would just be too much relative to the utility of the data (we would need like 250 shards) - retention, long term, we want to store data in Pulsar with up to unlimited retention so we can rebuild views. There is still some complexity there (like getting parallel access to segments in a topic for batch processing) but it is much further along than any other options - client APIs for consumer. We are a polyglot shop and really the only language where consuming Kinesis isn't terrible is Java (and other jvm languages). For every other language, we use lambda and while lambda is great it is still distinct deploy and management process from the rest of the app. Being able to deploy a simple consumer just as part of the app is really nice
- nostrebored 7y agoIs there any reason you've decided not to use Dynamo to manage your materialized views? Do you need to be able to generate a materialized view for a specific time window? It feels weird to me to use your pub sub system to handle your persistent storage for views, but I am definitely missing context into pulsar and your use case
- ryeguy 7y ago
- DevKoala 7y agoWas NATS a consideration for your use cases? At work, we are currently standardizing on NATS as our messaging system, and I would like to know if there is a valid comparison.
- liquidgecka 7y agoNats is not a replacement for pulsar or rabbitmq. It is a message passing system designed to pass lots of messages live, however if nobody is their to receive them they are lost and gone forever. There is a streaming layer but that is closer to Kafka and still does not provide the typical message model with an ack/nack API. I have used nats in several different ways but since it can be lossy its never been considered as a replacement for a pub sub message queue on my end. We used it for a chat message layer and that worked pretty well. As for its message passing layer that can be interesting but you end up writing all the retry and failure logic anyway so its usually just better to use an existing message layer that handles all of that for you anyway without all the funky abstractions. Again, its interesting but nowhere near close to being a rabbit or pulsar replacement if reliability is a goal.
- DevKoala 7y agoThank you for the observation.
- bauerd 7y agoSounds like zeromq to me
- JensRantil 7y agoSemi-correct observation. You could think of it as a ZeroMQ broker that supports the pub/sub patternz, yes. Difference is it's centralized (with support for HA), adds security, gives you monitoring, supports multi-tenancy and horizontal scaling allowing to connect multiple clusters. Benefit of using broker is that service discovery is a no-op. Also, NATS has client libraries for many languages that adds request/response semantics on top of the pub/sub semantics.
- staticassertion 7y agoWhat is the SQS-based system you are migrating from? I'm currently building a data processing system that is backed by S3 -> SQS based events, for persistent message passing.
- addisonj 7y agoWe have a number of systems, some use SQS, some use Kinesis. Part of the draw of Pulsar is having one piece of tech that we can unify everything over and offer more baseline features, like infinite retention via storage offloading or Pulsar IO connectors that standardize common operations. We aren't really targeting one use case, instead, we looked for the system that offered a broad set of features that other developers in the company want and is operationally doable with just a few people.
- ignoramous 7y agoOn behalf of everyone here, thanks a lot for answering every single question being asked. Highly appreciate it. I have questions myself: 1. Did it reduce (TCO) costs or increase it versus using Kinesis and SQS/SNS? 1a. Interestingly, there's no global-replication with those AWS services. Why did you require global-replication with the move to Apache Pulsar? 2. Since you mention internal auth: Weren't Cognito / KMS / Secrets Manager up to the job? Given these are integrated out-of-the-box with EC2? 3. Was it ever under-consideration to roll out pub/sub on top of Aurora for Postgres with Global Replication? https://layerci.com/blog/postgres-is-the-answer/ https://layerci.com/blog/postgres-is-the-answer/ Thanks again.
- addisonj 7y ago1. On a short time horizon, not as sure, back of the napkin, it took ~12 dev months (5 months with 2.5 people average on it). However, our cost per 1000 msgs/sec is much lower (like 1/4 the cost of Kinesis) so we fully expect that investment to pay off over time assuming that adoption by the rest of the org continues and we don't find a ton of issues. 1a. You are correct we didn't require geo-replication for existing use cases, however, initially, we saw geo-replication as an easy way to improve DR and we have an internal requirement for a DR zone in another region. Now that we have done the work, we are starting to see multiple places where we can simplify some things with geo-replication, so we think long term the feature will be really valuable 2. We split up auth into two main components: auth of users (where we use Okta) and auth of services. For okta, we just wrote a small webapp that users can log into via OKta and generate credentials. For apps/services, we already had hashicorp in place and wanted to just piggyback of our existing form of identity (IAM roles). Essentially, a user just associates an IAM role with a pulsar role and we generate and drop off credentials into a per-role unique shared location in vault that any IAM role can access (across multiple AWS accounts) 3. Once again, geo-replication wasn't really a hard requirement initially but more of something that we really like now that we have. I think the biggest reason why not postgres is that we have combined message rates (not everything is migrated yet) on the order of 300k msgs/sec across a few dozen services. Pulsar is designed to scale horizontally and also has really great organizational primitives as well as an ecosystem of tools. While I think you could maybe figure that out with some PG solution, having something purpose built really can pay big dividends for when you are trying to make a solution that can easily integrate into a complex ecosystems of many teams and many different apps/use cases
- unethical_ban 7y agoI have no idea what "pub-sub" is used for outside of its academic definitions, and I have no idea how Hashicorp Vault works - Don't you need a secret/password in cleartext at some point, for a given service or definition? You don't have to answer my questions, I am just shouting into the void. I'm glad it works for y'all.
- SlowRobotAhead 7y agoPub/sub messaging is super common. Most IOT devices that aren’t running HTTP stacks are using MQTT.
- NicoJuicy 7y agoEvent based systems -> which almost every microservice is nowadays. Or chatapps Or for IOT
- zapdrive 7y agoWe are developing a social app with features such as messaging, notifications etc. We decided to use Yedis [0] (Yugabyte Redis) which is a distributed Redis with persistence backed by RocksDB. Yugabyte supports multiple datacentre distribution. Yedis's pub/sub is distributed as well. We are already running a Yugabyte cluster for data storage in Cassandra. So we didn't have to do anything extra to get our distributed pub/sub up and running. Would you recommend using Pulsar instead? 0: https://docs.yugabyte.com/latest/yedis/ https://docs.yugabyte.com/latest/yedis/
- tschellenbach 7y agoNormally don't plug my own work, but this is super related. Did you ever check out Stream? https://getstream.io/ https://getstream.io/ We power chat and feeds for >500 million end users. Tech is Go, RocksDB & Raft based.
- zapdrive 7y agoYes, we did indeed consider Stream, but figured we could save some money by deploying and running our own system. We are very hopeful to quickly get a couple million users in a short time from our launch and that would have ran up our costs with stream quickly.
- tschellenbach 7y agoThat's nice to hear. Best of luck with your project. We've had some really large companies move to Stream from their in-house tech and save 30-70% comparing Stream's monthly fees vs their in-house hosting. (the difference gets much larger if you add the engineering & maintenance cost of their in-house systems). If your team is in the USA & funded it's pretty difficult to build in-house with a good ROI.
- addisonj 7y agoI am not familiar enough with either Yedis or your use case to make a recommendation, but I can say that Pulsar has a great set of features, particularly if you need long term retention, that make it very attractive. I also been impressed with the community and development pace. Being that the project is a top level Apache project and also has some adoption by quite a few different companies and a number of corporate sponsors the future of the project is pretty safe bet.
- oars 7y agoThat's amazing, thank you for sharing. I understand why you chose Pulsar over RabbitMQ, but wouldn't have Kafka been a good choice as well?
- biggestlou 7y agoThere are numerous places in the discussion where reasons for not choosing Kafka are elaborated.
- deleted 7y ago[deleted]
- jupp0r 7y agoThanks for answering so many questions on this! One more from me: Did you consider Google Cloud PubSub [1]? In general I'd be interested in your rationale for moving from a managed solution to something you maintain yourself, because in my experience the down the line costs of maintaining your own solution are often underestimated. [1] https://cloud.google.com/pubsub/docs/overview https://cloud.google.com/pubsub/docs/overview