Goodreads offloads DynamoDB tables to S3 and queries them with Athena
aws.amazon.com
aws.amazon.com
The main thing I've noticed in Goodreads since they were bought is a strong focus on Kindle integration. The newer Kindles all have Goodreads integration, you can send notes etc. to the Goodreads profile, the site itself has changed very little
(although Goodreads changed a lot of rules around the ability to delete user reviews if they focus too much on author behaviour around the time the Amazon purchase went through https://www.washingtonpost.com/blogs/compost/wp/2013/09/23/i... )
Like Goodreads, we have a fairly constrained catalog and we probably get a similar amount of queries.
The search results are now much, much better, but ES has a pretty poor edit distance / fuzziness algorithm so they still aren’t perfect.
I am about to start using ES for a project and just knowing which kinds of things could be useful to tune would be helpful.
Thanks!
I'll see if we can make a blog post, but here are a bunch of things that jump out (I work at a company that sells tickets to live events, so our users search for 'live events' like sporting events / concerts and 'performers' like teams and musicians):
- word order matters _a lot_. we did a lot of fiddling with the n-gram tokenzier (https://www.elastic.co/guide/en/elasticsearch/reference/curr...).. we ended up making word order matter a good amount (e.g., 'new york' vs 'york new' return very different results... considering them the same resulted in a lot of noise)
- where the user is searching from is pretty important -- we would fetch the 25 best results and then boost (i.e., reorder) them based on the user's distance from the event venue or the sports team's home venue. we also experimented with fetching more and more results (up to 250) and then boosting from this larger result set. note that ES couldn't take location into account out of the box -- we had to manually boost on the ES output
- we set up versioning with our autocomplete endpoint so we could more easily A/B test variants (highly recommend this)
- we built a system so non-technical employees could create "synonyms." for example, "nyc" could expand to "New York City." we also worked with our data science team to get a list of bad queries that might need synonyms to improve them. (we also automatically triggered a real-time re-index on synonym creation)
- we similarly had an "expectations" tool for bug reporting and finding patterns from common bugs
- we had to add a bunch of other metadata / suffixes to our documents. for example, we might want to return a 1pm Yankees game on August 4 when someone queries "august yankees afternoon game". so we have to interpret the time and add the month to what's being queried. similarly, we want this event to return when someone queries 'nyc baseball', so we need to ensure the league/sport is associated with the event document
- we also had to add "stop words" that we ignored when querying. these include 'game(s)', 'versus', 'concert(s)', 'tickets', etc
- we have an internal definition of performer or event "popularity", and needed to normalize this so ES's "match score" made more sense. (we had limited success here)
- their documentation describes fuzziness as: `fuzziness is interpreted as a Levenshtein Edit Distance — the number of one character changes that need to be made to one string to make it the same as another string` which is overly simplistic and really messy to override (we decided against it)
- because we had two different entities in our results ('events' and 'performers'), we had to figure out how to compare different entities (it was generally easier to compare results within entities) based on what was returned, time to event, location of event, and home location of the performer. we also added additional entities / pages on an ad-hoc basis which further complicated things
- we also needed to exclude low quality performers and events from our catalog (e.g., performers with no events, events with no tickets for sale)
In addition to configuring ES, it was pretty difficult to settle on a KPI because it's not that easy to put searches in the context of the entire user session... we could see if a given query resulted in: the user clicking on a result, or no search results, or the user deleting everything in the box and starting over, but we had a hard time following the user and seeing if the click led to a purchase.
Also, as a disclaimer, I didn't actually write any code for this project (I'm a product manager). But I did take a computational linguistics class in college and worked very closely with the developer :)
I'm hoping you do find the time to write a blog post on this.
Thanks a lot!
(Disclaimer: I used to work with one of the authors at Eventbrite)
IIRC you originally could add your own books to Goodreads which let you fill in the gaps.
[1] https://www.goodreads.com/book/show/1869210.La_insidiosa_fat...
[2] https://www.goodreads.com/book/show/101558.Harry_Potter_y_el...
That would mostly create a good base of about 90% of my books in the last 10 years.
The DB scrape uses a template from Data Pipeline that under the hood uses the Dynamo DB Scan API. Not really surprising, but that API uses JSON. I wanted to use as much off the shelf software as I could to get data into Athena.
In this architecture Lambda is only used to listen to the SNS topic that fires when the Data Pipeline job succeeds or fails so we’re not pushing the limits of Lambda at all. You’d probably hit an EC2 limit Witt Data Pipeline before hitting your Lambda limit on the account.
https://twitter.com/iamsomewalrus/status/1016368717880389632
As mentioned in another comment I’ve found having Dynamo snapshots in Athena really useful as an oncall to sanity check snapshots (what was the state of Harry Potter 3 months ago compared to now?) and to answer product questions that can only be answered from the raw production data.
Motivating example: you have huge tables in Redshift that are either infrequently accessed or the usefulness of the data decays over time (website logs, customer order information). In this scenario you're paying a lot just to keep data in Redshift (storage) but a large subset of the data is laying dormant (no compute).
If you're bought into the Redshift ecosystem this is where Redshift Spectrum comes in. If you're a smaller company you could just store the data in S3 and "spin up" the compute when you need it (Athena, Glue jobs, or Elastic Map Reduce clusters).
(Disclaimer: I work for Google)
I would never pick Google for anything important.
I've been exploring the former and it seems to only make sense if the size of your data is at a scale that is beyond what a single SQL database instance can handle, and even then, you can continue to scale out with systems like Citus so the limit isn't a hard one. SQL gives one so much (data mutability, consistency, indexes, etc.) that I am hesitant to give it up unless the tradeoffs make sense.
However, it is also very useful if the following two things are true: 1. You have a very large stream of incoming structured data that is mostly write-once-read-never, like logs. 2. Your query use cases are relatively simple and static. If those fit your use case, then S3 + parquet + Athena is very easy and very cheap.
Goodreads hit scaling issues a while ago with Active Record and a single database so we broke up the data into separate MySQL servers. At that point joining data across DB servers is impossible so we went with Redshift for BI. Nowadays we would probably go with a datalake on S3.
I think it makes sense when there's no in-place updates; either querying write-once data like logs or the output of batch data processing roll-ups that replace the previous data. The less you need the relational model (like joins), the better, but some of those needs can be met through careful design of the storage schema and denormalization.
I wouldn't advocate this sort of solution if your requirements include in-place updates of existing data, frequent/granular updates of new data, expressive ad-hoc queries that use the full capability of relational algebra, or tight latency requirements. You also lose the safety net of referential integrity and table-level constraints, as those are now enforced in custom code that can have bugs.
I would say maintaining this system cost about a half-engineer for ongoing maintenance and new functionality.
S3 + SQL is good for huge log/machine data, exploratory use cases that are not yet productionized, ELT (to get data from raw files into SQL, used as a feed to later layers), quick and dirty SQL against a directory of similarly structured files. I tend to think of it as a utility layer.
For long term analytics use, that involves a domain model, I’d still stick with dimensionally modeled (or snowflake) data warehouse techniques. Getting data into such a model can take weeks to months, so sometimes it might be better to do something quick and dirty in a data lake to prove a dataset or get a quick answer, vs. slow down the business waiting for a perfect model.
Lastly, I see storage + SQL as being the same conceptually as any RDBMS, with different performance, cost, and functionality. For example, SQL Server proprietary disk format + SQL Server query engine is somewhat analogous to Parquet + PrestoDB. In fact many proprietary vendors integrate with HDFS as a distributed storage layer for their proprietary formats which can be queried alongside open source storage formats by proprietary SQL query engines too.
I think it compares more with something like BigQuery but if you already have your data in S3 maybe you get a more well integrated system if you stick with AWS tools.
* The default timeout is ~ 48 hours and you pay per Data Processing Unit (DPU) that you've provisioned the Job. * Currently it supports Python and Scala. As far as I'm aware you can't run Java jobs directly, but you can upload JAR libraries and use them in your code.
Re serverless vs dedicated / ephemeral clusters: Like with any serverless runtime environment you are trading convenience (across a few dimensions) for flexibility.
The Glue environment runs in a few limited runtimes and uses a specific version of Spark that you have no control over updating. Given that it's pretty quick to author a job, you can set the required DPU and Glue handles that, and you don't have to worry about sizing the cluster for your data size. For me most of my jobs fit within those constraints.
At some point on the cost curve it may make sense for you to move all of your jobs from Glue into a dedicated cluster on EMR. You may also get there sooner if you need to use specific frameworks or libraries.
I haven't had a chance to actually read your post yet, but I'm really looking forward to it.
We used Terraform instead of CloudFormation, although there are a few places it didn't/doesn't cover.
We're also doing protobuf -> parquet instead of JSON; Scala instead of Python, and one of our major feature/issues is that the incoming data is out-of-order; Firehose outputs to partitions based on when the event arrived, and we're repartitioning based on when the event occurred (according to a timestamp in the message)
I see you ran into the capitalization issue too :) I ran across docs somewhere that said something along the lines of Glue downcasing column names, which definitely fits observed behavior.
I'm hoping to turn what we've done into a similar post, although I think it'll be a month or two before I can get to that.
Edit: Oh, I also got a lot of mileage out of the Zeppelin Notebooks. Way better than a raw dev endpoint, but watch out on the cost for both :)
The notebooks can be provisioned with almost a single click from the Dev endpoint console, but they changed the recipe halfway through my work on it, and now you have to also SSH into the box and run a script to setup some of the security. :/ Still totally worth it, tho.
Thumbs up on the Glue Dev endpoint. It's been killer. I had trouble setting up a Notebook (I wanted to get fancy with Docker) and I usually use the Python repl link that's provided.
I'm working on a follow up post that removes the Data Pipeline -> Lambda and uses the new Glue DynamoDB integration.
Looking forward to your blog post!
What was the trouble you had with a notebook? I can probably post up some of our Terraform code (which includes notes on the parts Terraform doesn't cover).
Oh, yeah - there was also the S3... VPC? endpoint. That needed to exist.
There were a lot of wires, and Amazon documentation is decent as a reference but rubbish as a tutorial :/
[1]: https://aws.amazon.com/blogs/aws/amazon-s3-performance-tips-...
The first major project after being acquired was to make a pared down Goodreads experience available on the Kindle Paperwhite. Our first services came out of that initiative to provide a buffer between the Kindle traffic and the Goodreads Rails app.
That being said I’ll be the first to caution small teams should avoid microservices at first for fear of creating a distributed monolith.
What is your heuristic for determining when it is the right time to look at microservices?
Engineering team size: you have teams large enough (~3-4) to own and iterate on a subset of related functionality for the long term.
Tooling: you have a builder tools type team that provides a tooling and observability happy path.
Traffic scale: you have functionality that operates at 2 or more order of magnitude higher that the rest of your application.
Decoupled: you have functionality that can be decoupled from your main app and, most importantly, isn’t required for your main app’s uptime. Like a search service or something.
A combination of the above.
Reach out to me feeneyj @ amazon if you'd like I can forward you to the right people.
They're cross-charged with lots of internal book-keeping. The rates teams "pay" internally are different from public pricing of course but expensive things externally are still expensive things internally. The cost-accounting used to be (3+ years ago) much more "just looking...no pressure to keep them down", but recently there's been huge efforts to bring internal AWS costs down, especially for EC2 usage. I heard rumors that much of the Prime Day fiasco this year could have been avoided if teams would have been permitted to spin up enough capacity.
Meanwhile he was deploying some machine-learning categorization stuff he was working on to some cluster that would cost five figures of compute each time, and nobody batted an eye.
Also you have to jump through pretty extreme hoops if your hardware estimates (usually made at least 3 months out) were under-shot and now your service is redlining. This leads to teams way over-estimating their hardware needs and thus millions being spent on idle reserved EC2 capacity. So now they police (with savagery) idle capacity, so services or workloads that are "bursty" are basically an internal-bureaucracy nightmare.
It's much nicer to use AWS outside of Amazon :)
Just BI.
As I mentioned in another comment I, personally, wouldn't advocate for new teams or startups to use a service oriented architecture right out the gate. It's too easy to end up with a distributed monolith (circular dependencies between services, services uptimes that are tightly coupled). Engineers also tend to underestimate the build tooling, observability excellence and discipline you need to make it seamless.
We’re trying to get to a place where we have the data in S3 for engineers to build products off of and for the oncalls to do sanity checks and the data in Redshift for our BI needs.
The other replies got it right w.r.t. other BI tools. If you’re using Tableau I think it integrates with Redshift, right? In that case Redshift Spectrum is an option.
If you don’t have any existing BI tools then Quicksight is an option or alternatively you can spin up an Elastic Map Reduce (EMR) cluster with your fav open source BI tools
However, I read a harrowing / awe-inspring blog post about someone doing just that. So...¯\_(ツ)_/¯
https://docs.aws.amazon.com/amazondynamodb/latest/developerg...
Now I've learned to read the Best Practices section of any AWS service before I start implementing. It saves a lot of heartache
There are no indexes in the traditional MySQL / Postgres sense of the word. You can, however, layout the data to make your querying more efficient. See: https://docs.aws.amazon.com/athena/latest/ug/partitions.html
From what I can gather, Amazon acquired them and now they have to figure out something to do. There is plenty of UX to fix.
Joe feeneys HN account is indeed iamawalrus; he wasn’t impersonating anyone. Further his comment just said:
“Hey! I’m the author of this post. I’m pretty chuffed to see this here. Happy to answer any questions.”
I want to hear what he has to say.
Looks like some automated filter, for sure.
Wouldn't using ETL Glue job to directly dump data in parquet format be better?