The denormalized query engine design pattern by Simon WIllison

This video features Simon Willison at DjangoCon US 2017 in Spokane, Washington, USA.

The denormalized query engine design pattern by Simon WIllison
0:43:11
Published September 8, 2017
2,549 views

DjangoCon US 2017 - The denormalized query engine design pattern by Simon WIllison

Most web applications need to offer search functionality. Open source tools like Solr and Elasticsearch are a powerful option for building custom search engines… but it turns out they can be used for way more than just search.

By treating your search engine as a denormalization layer, you can use it to answer queries that would be too expensive to answer using your core relational database. Questions like “What are the top twenty tags used by my users from Spain?” or “What are the most common times of day for events to start?” or “Which articles contain addresses within 500 miles of Toronto?”.

With the denormalized query engine design pattern, modifications to relational data are published to a denormalized schema in Elasticsearch or Solr. Data queries can then be answered using either the relational database or the search engine, depending on the nature of the specific query. The search engine returns database IDs, which are inflated from the database before being displayed to a user - ensuring that users never see stale data even if the search engine is not 100% up to date with the latest changes. This opens up all kinds of new capabilities for slicing, dicing and exploring data.

In this talk, I’ll be illustrating this pattern by focusing on Elasticsearch - showing how it can be used with Django to bring new capabilities to your application. I’ll discuss the challenge of keeping data synchronized between a relational database and a search engine, and show examples of features that become much easier to build once you have this denormalization layer in place.

Use-cases I explore will include:

Finding interesting patterns in your data
Building a recommendation engine
Advanced geographical search and filtering
Reacting to recent user activity on your site
Analyzing a new large dataset using Elasticsearch and Kibana

This talk was presented at: https://2017.djangocon.us/talks/the-denormalized-query-engine-design-pattern/

LINKS:
Follow Carlos Martinez 👇
On Twitter: https://twitter.com/simonw
Official homepage: http://lanyrd.com/profile/simonw/
Github: https://github.com/simonw/

Follow DjangCon US 👇
https://twitter.com/djangocon

Follow DEFNA 👇
https://twitter.com/defnado
https://www.defna.org/

Summary

Simon Willison names and explains the “denormalized query engine” pattern: keep a relational database as the source of truth, copy query-oriented data into a horizontally scalable search engine such as Elasticsearch, and synchronize the two carefully. Search indexes handle counting, aggregations, multi-field queries, relevance, text, recommendations, and geographic queries more effectively than many relational queries, while smart query routing and fetching final records by ID from the database prevent users from seeing stale data. He compares synchronization through timestamps and polling, queues, and database replication logs, and describes practical techniques including cascading updates, self-repair for deleted records, and accurate filters that combine recent database changes with indexed results.

Key takeaways

  • A relational database can remain the authoritative store while a denormalized search index serves complex, high-volume queries.
  • Search engines are useful for much more than full-text search: they provide fast counts, nested aggregations, relevance scoring, recommendations, and geographic filtering.
  • Synchronization is the difficult part and can be implemented with timestamp polling, queues, or database replication logs feeding systems such as Kafka.
  • Routing a user’s own data to the database and public or other users’ data to the index helps hide indexing delays.
  • Returning IDs from the search engine and loading current objects from the database avoids serving stale indexed data.
  • Recent, authoritative database changes can be combined with indexed results to make filters appear immediately up to date.

Summarised automatically from the transcript.

Chapters

  1. 0:00 Career Background and the Design Pattern Simon Willison introduces his career path and names the denormalized query engine pattern.
  2. 1:46 Relational Databases and Search Indexes The talk explains how a relational database can remain the source of truth while a search index provides additional query capabilities.
  3. 3:18 Relational Database Limitations Counting, large scans, and multi-field queries expose weaknesses in conventional relational database indexing.
  4. 4:54 Search Engine Strengths Search engines provide horizontal scaling, fast counts, aggregations, relevance scoring, and full-text search.
  5. 6:26 Flickr’s Sharded Search Architecture Flickr combines sharded databases with a search index to handle queries spanning many shards.
  6. 11:06 Lanyard’s Search-Based Querying Lanyard replaces an expensive relational query with a Solr query over denormalized attendee data.
  7. 13:25 Elasticsearch Fundamentals Elasticsearch is introduced as a JSON-over-HTTP search and analytics engine built for horizontal scaling.
  8. 15:51 Aggregations and Faceted Navigation A congressional email explorer demonstrates interactive filtering, nested aggregations, and fast analytical queries.
  9. 23:00 Database-to-Index Synchronization The talk compares timestamp polling, queues, and database replication logs for keeping a search index current.
  10. 30:05 Fresh Results and Consistency Techniques Simon describes ID-based lookups, self-repair for deleted records, and the accurate-filter trick for avoiding stale results.

Transcript

8,565 words · auto-generated Show

Automatically transcribed, so expect mistakes in names and technical terms.

0:14

Speaker 1: So yeah, um, so I will start with a little bit of a career introduction, which I promise is is very relevant to the talk. I'll be talking about a design pattern that is sort of I've stalked throughout my career. So I started out many years ago at the Lawrence at a tiny little local newspaper in Kansas called the La Lawrence Journal World working on a web framework that eventually became Django about a year after I'd left there. I moved on to work at Yahoo where I was I briefly tinkered with the Flick team and then worked on various like product development and research projects. I did data journalism at the Guardian, which was the most fun job ever because it's you get to work with data on like journalism deadlines which you know that sort of ties back into the original tagline for Django as well. Then after the Guardian I did a startup.

0:59

Speaker 1: I co-founded Lanyard with um my my wife Natalie who's there in the front, uh ran that for three years and then sold that to Eventbrite. So now through a various path of of of different um of different machinations, I'm an engineering director at Eventbrite over in San Francisco. But the thing I want to talk about today is a design pattern. And design patterns, really the power of design patterns is almost entirely in the name. You know, no nobody really invents design patterns. You more sort of look at something that other people are doing and you slap a name on it, and then it becomes something which people can talk about. And the pattern I want to describe today is one which, to my surprise, no one else seems to have slapped a name on yet. So I'm slapping the name on it and I want to start getting discussions going because I think it's a pattern that can help out with a lot of different projects.

1:46

Speaker 1: um in in a lot of different ways. And the name I've picked for this pattern is the denormalized query engine, which I hope is just snappy enough that it'll work work for people. And essentially this is a way of working of of taking a system built on a relational database and enhancing it using a search index such that you can do a huge amount of additional um interesting things with it. So the key idea is you have your relational database as your single point of truth. And you know we we all got kind of infatuated with NoSQL a few years ago. I feel like that infatuation has worn off a little bit. It turns out that 40 years of computer science has made relational databases a particularly reliable place to keep the data that you care about. But anyway, you you have your data in your relational database.

2:32

Speaker 1: You then denormalize the relevant data into a separate search index. So you take all of that data, you think think about the bits that would make most sense to be um to be denormalized to be queryable in different ways, and you get those into a search index, and then you invest an enormous amount of effort in synchronization between the the two, making sure that whenever somebody changes something in that database, you get that into the search index as reliably and as quickly as possible. And that's the hard bit. I'll be talking a little bit more about some strategies For doing that towards the end of this talk. And why but why would you want to do this? Well, really, this is a way of addressing some of the weaknesses that most relational databases The first one which I'd imagine many people have run into is relational databases aren't really very good at counting things.

3:18

Speaker 1: If anyone's ever implemented pagination where you have like two hundred thousand rows in a table and you want to do page one, page two, page three, you'll find that the bit where you count select, uh where you count star against that table is the bit that actually starts to hurt you first, because the database has to scan through all 200,000 rows just to generate that. account. As a general rule, anytime you're doing something that end users are um are going to be accessing, you need to avoid queries that read more than say a few thousand rows at a time. Relational databases are insanely fast at primary key lookups, and they're insanely fast at range queries against an index. But if you've got a query that needs to query needs to invest needs to look at 10,000 rows, that's going to add up to one, two, three seconds, and it's going to be something you can't deploy in

4:05

Speaker 1: an application that end users hit. all of the time. And there's one that is currently specific to MySQL, I believe Postgres has fixed this one. But MySQL can only use one one of the indexes defined in the data. database for a query that's being executed. So you might think that you can slap indexes on the age column and on the um and on the uh uh and on the the the job title column searches across both of those at once. But actually it'll pick one of those two indexes and it'll use that to speed up your query. So actually the moment you start doing more complicated lookups, the database indexing scheme really starts um starts undermining you. Meanwhile, search engines have a whole bunch of strengths. Firstly, modern search engines, and I'm mainly talking, I'll be mainly talking about Elasticsearch in this talk, but the same is true for solar and other search engines

4:54

Speaker 1: engines as well are really good at scaling horizontally. Like you you can take a system like Elasticsearch and literally just throw more um machines, throw more nodes at it, and it will rebalance across that that full cluster and give you more read performance, more write performance. performance and and just jet general um um and and improvements in your capacity as as you scale that up. Um they're really good at counting. Um databases not It's not so great at counting. Search engines are really, really fast at this, and they're great at aggregations as well, which I'll talk about in more detail. a moment. You can run queries across multiple indexed fields. So if you have a super complicated query where you're you're um you're looking at like four or five different fields And combining those together, a search engine will make short work of that.

5:39

Speaker 1: They're unsurprisingly very good at relevance calculations and scoring because that's kind of the nature of of uh of the beast. And they give you text search. You get all of these benefits and you can implement full text search as well. I deliberately left that one to last because um my interest in search engines goes way beyond just using them to search for text that users have entered in. I think that this entire design pattern revolves around the fact that search engines have strengths beyond just being able to implement a full Text search. So I'm going to roll back in time to 2005 to talk about the first time I saw this pattern in the wild, and that was at uh Flickr, the photo photo sharing site. So back in 2005, Flickr were having enormous scaling problems

6:26

Speaker 1: because it was the birth of Web 2. 0, it was social meets it was, I don't think, social media. Media was even a turn back then. They had an enormous quantity of users coming into the um adding to the service and uploading photos. And as an industry, we hadn't really figured out how to do this web-scale engineering thing yet. Sites like Flickr were having to figure this stuff stuff out from scratch and figure out how to scale up to handle these these giant um giant numbers these enormous numbers of users and huge amounts of data. The CTO at Flickr, uh Cal Henderson, um, wrote a book about this called Building Scalable Websites, which came out over a decade ago. now and I think is still very relevant today because it essentially talks through the lessons they learned at Flickr figuring out how to how to scale these things up, how to how to build a scalable um

7:12

Speaker 1: web application. And the technique that they used at Flickr that got them out of their hole was was uh database sharding. Um so this is a very common technique. um to this day. Essentially what you do is you say, okay, we had one MySQL database and we couldn't keep up. There were just too many writes coming into this database. So what we'll do is we'll split it into multiple databases and we'll put different users on different shards. So maybe we'll put users one through ten thousand on this database, 10,000 to 20,000 on this database, and so on and so forth. And this is a very naive description of sharding, but I hope it illustrates the The concept. So if you do this, life becomes a lot easier because you can as your user base grows, you just add more hardware, you add more databases. And if you want to do things like

7:58

Speaker 1: show me the most recent photographs uploaded by Simon, you say, okay, well Simon's on So I'll go to SAD 3 and I'll select photos from there ordered by whatever. And that'll give me an answer to my question. So that sounds well, I hope It sounds relatively straightforward, but there's one massive problem, which is what you do with data that hap that that lives across multiple different shards. A great example at Flickr. Flickr um Flickr were very early adopters of the idea of sort of user-provided tags. And so you can go to Flickr today and you can see all of the photos that have been tagged racco tagged with the raccoons tag and see though see the most recent uploads and all of that kind of stuff. stuff. And the obvious problem here is if you've got photos across five or six or a dozen or a hundred

8:44

Speaker 1: different sharded databases and you need to find all of the raccoon phot all of the photos tagged raccoons, are you going to do a query that hits a hundred databases at once and then try and combine the results as they come back. That's not really a sort of practical way of solving this problem. So what the Flickr team did is they took advantage of the fact that they were now within within Yahoo and they leaned on a piece of Yahoo internal technology Technology called Vespa, which was pretty much what Elasticsearch is today, but 10 years ago and written in C and kind of gnarly to work with. And so what they did is they said, okay, we're gonna have our sharded database. We'll have different users' photos, we'll differ live in different places. That's all fine. And then we'll have a search index which we load all of the photos from all of the shards into. And this was on Vespa, which could scale horizontally

9:31

Speaker 1: and gave them all of those benefits. And then we can um when somebody tr makes a query against Flickr, we can make a decision. We can say, if it's you and you're looking at your own photos, that's going to be a database query. If it's you looking at other people's photos or if it's you looking at every public photo tagged raccoons will turn that into a search query instead. This is a diagram from Aaron Strap Cope, one of the engineers at Flickr who worked on this at the time. And this turns out to work really, really well. There's a key concept embedded in here, which is the this idea of smart query routing. Because if you're building software that human beings use, one of the worst things you can do is have a bug where somebody makes an edit and then they refresh the page or they go to or the uh the the UI refreshes and their edit hasn't isn't shown

10:17

Speaker 1: back to them. If a user tags a photo raccoons and then goes and looks at their photos tagged raccoons and it's not in there, this is a bug and they're they're they're justifiably annoyed by it. So the the the solution flickr um the solution Flickr used was to say if you're looking at your own data, that should be a relational database hit because we can't guarantee that the search index has got those changes yet. Generally with these systems it can can take a few seconds up to a few minutes for the underlying search index to reflect those changes. If you're looking at other people's data, you're just not gonna you you you're never gonna know if your friend just uploaded photo tag raccoons like five seconds ago and you can't see it yet, that's not something you'll be able to observe. So at that point it's safe for us to use use the search index for public and other people's data and the relational database for our own data.

11:06

Speaker 1: So the next uh um well fast forward a few years to 2010 um when we used as we used Solar, um another open source search engine, to solve a similar kind of problem. So we launched Lanyard in my wife and I launched Lanyard on our honeymoon. It was supposed to be a side project and it ended up growing way, way beyond that. But the um the the idea was we we were we were in uh Casablanca in Morocco and we had food poisoning and were unable to keep on travelling and um Casablanca it was during Ramadan when none of the restaurants were open So we couldn't get anything to eat anywhere else either. So we figured, okay, we'll rent an apartment for two weeks, we'll look after ourselves, we'll we'll cook ourselves better, and we'll try and ship this side project that we'd been working on.

11:52

Speaker 1: We made the mistake of building a side project with with user account. logins, which you should never do because users have expectations. And we also built it on top of Twitter. And so the core feature of Lanyard when we launched was you sign in with your Twitter account. And we show you conferences that your Twitter friends, the people you follow on Twitter, are speaking at or attending. Which When you think about that in terms of a database query, ends up being a SQL query where you say, select star from events where it's in the future and at least one of the attendees is in this list of a thousand IDs that I can I've pulled back from Twitter. This is the kind of thing that relational databases are incredibly bad at. And so we we launched, um, we got a flurry of initial activity, we ended up on TechCrunch UK unexpectedly, and the site just died instantly

12:38

Speaker 1: Because the most popular page on our site was the page with the most expensive database query. And so the way we resolved this was we'd started using Solar to provide search and we thought, well, maybe this is something we can redefine as a search problem. So we turned this feature from a giant hairy SQL query into a solar search where we said, hey solar, I want events that are in the future where the attendee IDs field in solar matches any one of this list of applications. to 2,000 um IDs that I've pulled back from Twitter. And do that and order by date. We built this thinking this will probably do us for like a few months. until we can figure out a better solution. And now five years later that's still how this page works because it turns out search engines are incredibly good at exactly that kind of query.

13:25

Speaker 1: Normally a search engine expects to be dealing with words because users type words in but actually things like user IDs or other terms work ex exactly as well and the underlying architecture of the search engine is really good at um dividing those up search uh at merging together all of the documents in the index that have fields that match whatever criteria it is you're passing in. So this is a good time to switch over to talking a little bit about Elasticsearch. As I mentioned, Lanyard was so Flickr used Vesper, which I looked this morning and it turns out Yahoo have actually out open sourced it, which is kind of kind of interesting but um I don't know if it's it's gained an enormous uh amount of community adoption yet. Um we used solar at Lanyard which is a very fine search engine albeit one which was clearly a

14:13

Speaker 1: It was clearly a um, it's a sort of product of the time in which it was designed. It's very fo it's very XML heavy. They sort of added JSON as a as a as a as a later uh later detail. Elasticsearch is what you would get if you designed solar today, if you said, okay, clearly the world speaks JSON and the world speaks HTTP. and we're going to want things to scale horizontally. Let's combine all of those things together and build a really good search engine. And so it's an open source search engine. It's built on top of the same Lucene search library that Solo and other projects are used in the past. The interface is entirely JSON over HTTP, which is great because it means you can talk to it from any programming language that speaks those two things, which I'm pretty sure is everything these days. And it makes it very easy to use from for, you can use it as a sort of

15:00

Speaker 1: central point between different programming languages very easily as well. The marketing bump all claims to be a real-time search engine. In practice, it's close enough. We're talking a few seconds between you submitting a document into Elasticsearch and that document becoming available across the cluster for you to run Queries against. It has an insanely powerful query language. I'll show you a little bit of that as we go along, but essentially there's a domain-specific language written in JSON that lets you construct extremely powerful and complicated search queries. And it also has a very strong focus on analytics. If you go to the Elasticsearch website, you'll see that they mainly talk about it as a analytics engine for things Like log analysis and so forth. That's where a lot of their focus has been. And as I said earlier, this is the thing that excites me about search engines, is sure, full-text search is nice, but being able to do these more complicated analytical

15:51

Speaker 1: questions Is where they get really interesting. And then finally, the elastic really does mean elastic. Elasticsearch scales horizontally. I've actually, um in the past, I've run a cluster of four nodes and just killed one of them at random to see what would happen. happen and you can watch the documents rebalancing across the remaining loans in real time using various some visualization plugins. So it's it does live up to the to the E in its name. So I've talked a little bit about how I'm excited about more than just search. The feature that I'm specifically excited about is Aggregations and the best way to illustrate those is with an example. So this is a project that I built last year. It's a little side project. Embarrassingly, since I'm speaking at JengaCon This is actually the only thing I've ever written in Flask, because I decided to try out

16:39

Speaker 1: see what Flask looked like. Flask worked very well. It was a very nice way of building this. So what this is, is um It's called DC Inbox Explorer. And there is a project at um Stevens University. Is it Stevens? Which called DC Inbox, which collects the emails that senators And Congresspeople send out to their constituents. So this um this researcher subscribes to all of these different mailing lists and gathers all of those emails and puts them in a giant giant um a giant JSON file that you can then Then use to run research about who's emailing about what and when. And because it was a giant JSON file, it was very easy for me to take this and import it into Elasticsearch. The source code for this is all available on GitHub and there's not very much of it.

17:26

Speaker 1: If you want to see how this works and get an example of Elasticsearch, this is a pretty good starting point. So basically what this does is it shows you all 57,000 emails That have been collected by this project. And it lets you search. So I can search for say security. And then I get back 15,000 results. It shows um the number of emails sent by month at the top. And then down the side, this is the this is my sort of pet favorite feature of any piece of uh of search software. It has these things which are sometimes called facets, sometimes called filters. So what this is saying is that without a search term, there are 56,000 emails Emails, 36, 7,000 of those were sent by Republicans, 19,000 were sent by Democrats, 280 were sent by Independents. You can see them broken down by by representative versus senators, by the state that the um that politician represents.

18:18

Speaker 1: And if I then search for, say, security, those numbers update. So now I can see that the emails that mention the term security 2,500 of those were sent by Senator. If I click on that, I'm now seeing emails sent by a senator that mention security. And I can see that of those, 810 were Democrats. um 1600 Republicans, I can now see that the state that is most concerned about the state whose senators care the most about security is a practically M. Is that Maine? And I get this, and so I can keep on drilling down and say, okay, set emails mentioning security from Maine by I've got gender in there, I can look at by by male or female senators. I can see the actual sentence themselves. So Susan Collins is the most prolific emailer from the state of Maine on the subject of security.

19:05

Speaker 1: But the key thing to illustrate I'm illustrating here is that when you've got a search engine like this elastic search these kinds of calculations these um aggregate counts become uh essentially free for you to that they're very they they become very easy from an implementation point of view the performance is super is is is super fast. So you can build this kind of highly interactive interface very easily on top of that underlying engine. And if you go on sites like Amazon and Book Booking. com and so forth, they all make very extensive use of this faceted navigation pattern. If you try to do this with a relational database, you're likely to run into some pretty nasty performance problems. problems pretty quickly. What I can also show you is under the hood, I'll show you a little bit of what Elasticsearch itself looks like to work with. This is a tool called Sense, which is kind of an IDE for talking to Elasticsearch.

19:55

Speaker 1: And as I mentioned, Elasticsearch is all JSON and HTTP. So I'm going to do a get against the slash email slash underscore search endpoint. And this returns essentially all of the emails it's paginated. I think it returns 10 to a page, and you can see at the very top it says that there are 58 that 56,000 of them. This is what an email looks like. It's got all of this data that I ingested. when I indexed the documents. And so then I can say, actually, you know, I want to run a search. So this is illustrating the Elasticsearch domain-specific language. I'm saying I'm doing a get against email slash search. I want the I'm I'm it's a query and I want to match the term security in the body field. And I'll run that. And now it gives me back 15,000 emails, and those are the ones that that match this particular query.

20:43

Speaker 1: There's a slight oddity of Elasticsearch. This is an HTTP get, which includes a body as if it were was an HTTP post. I had no idea this was even possible, but apparently it is. And um Elasticsearch uses it for everything. So um so that that that's something I I learned in playing around with this. But let's go a step further and say, okay, we're going to search for all of the emails matching security, but I want to also get numbers broken down by role type, which which is senator versus representative, and by party. And if I run this search here, I get back all of these different search results. And then at the bottom I get this aggregations block where it says role type representative has 12,500, senator has two and a half thousand, Republicans, 10,000, Democrat, 5,000, Independent. And it's 123.

21:29

Speaker 1: These are the numbers that you saw in the interface earlier, but this illustrates how it's just a little bit of extra JSON that you add to your query. And the um query time for this was four milliseconds, um, which I think That's pretty good, you know. I'll show one last example just to illustrate something that I think is unique to Elasticsearch, which is that Elasticsearch lets you take these aggregations and nest them. So here what I'm doing is I'm saying I want to aggregate counts by the party, and then within that party I'd like to like to do counts by the role type. So if I run this, I get back results where you can see that the Republicans have sent 37,000 emails. Of those, 31,000 were sent by representatives 5,000 sent by senators. The Democrats, 19,000 emails, of which 15,000 were representatives and 4,000 were senators.

22:14

Speaker 1: And then the independents at the bottom, there's apparently one email sent by an independent representative, and the rest were all sent by senators. I'm intrigued. Let's have a look. So if I do independent here and then representatives, sure enough, uh representative Gregorio Sablan. um has sent a single email apologizing for spam, which cons considering he's only sent one email is a little bit surprising. So there we go. So um that's sort of some of the power that you get once you start adding an an uh an engine like Elasticsearch into your Stack.

23:00

Speaker 1: I said I'd earlier I'd talk about the difficult problem, which is the synchronization strategy between that and your relational database. And I've tried a whole bunch of different ways of doing this. The three that I've had the most luck with are the three um so I'll talk through these in a little bit more detail. But to um to sort of to to um to repeat the problem we're trying to solve here is you've got users who are making changes in your relational database. They're updating things, adding things. You want those to be reflected in your search index as quickly as possible because any delay could result in like strange behavior that your users don't understand. And you want to do this in a way that is performant and efficient and doesn't cause too much overhead on the on the various parts of your stack. So the simplest way to do this is to basically keep it in the database and it's to have a last touched or a changed

23:49

Speaker 1: touch. timestamp in the um in the on the actual rows of your database, which gets updated any time somebody changes that row. This is a very common pattern. Here's what it might look like in the Django RM. I've got the last touched column, it's a datetime field, db underscore index equals true. It's important to stick an index on this because you're going to be polling this from a cron job like once a minute. So it needs to be able to return results quickly. And you set it to default to now. And that's fine. And then if once you've got this set up, the simplest thing to do is just have a cron that runs once a minute, selects star from um from that table where the uh where change state is within the last minute and then re-indexes those items. The nice thing about having this as a timestamp is that an indexer can keep track of the last time

24:36

Speaker 1: last time that it polled, the last thing it saw. So if your indexer doesn't run for five minutes, when it runs again it can catch up on five minutes. minutes worth of changes all at once. There is a subtlety to this, which is that quite often when you're building a search index from a relational database, there are changes that happen to other tables which still should trigger an update. So on Lanyard we have a concept called a guide where a guide is a somebody might create my guides to JavaScript conferences in Europe and they'll then add conference they'll they'll add events into that guide. The problem with that, and and we try and include the name of that guide in searches. So if you search for JavaScript Europe, as long as an event has been added to the JavaScript in Europe guide, it should show up up. What that means though is anytime somebody changes the guide, we need to re-index all of the events that that link to that guide because there's a bit of dependent data that has now been updated.

25:27

Speaker 1: And if you're using the last touch mechanism, that's pretty easy. You can say anytime a guide is edited, do guide. conferences. all, find all of the conferences attached, and update their last touch to date to the current timestamp as well. For the most part, that works fantastically well and this means that you get these sort of cascading changes happening within your database which your search index can then then track. catch up on. So a slightly more sophisticated way of doing this is with a queue. You can have application logic that says anytime somebody updates an event, write that or updates a document, write that document up ID into a queue and then have something at the other end of the queue which is consuming from it and re-indexing those documents. A really nice uh side effect of this is that you

26:13

Speaker 1: you can have deduping, although you get de-duping with the um the previous mechanism as well, but you can write your indexes so that it says, okay, there's been a flurry of activity around this particular document, but I'm gonna batch those up and a few seconds Later, I'll do one in one re-indexing call to recreate that in Elasticsearch. So I've built this a few times. I've built this on top of Redis, which worked great as a little Heroku app. Um I've built this on um at EventBrite, we use Kafka for this. And in fact, we we have a slightly more sophisticated system, which I'll dive into in a moment. Um another nice thing about cues is if you have a persistent queue, you get that replayability as well. So you can replay all of the indexing changes from the past five minutes. This is the most sophisticated way that

27:00

Speaker 1: I've seen this solved. And this is what we do at Eventbrite , which is tap into the database's replication log itself. So MySQL has very robust replication log It's very easy to have a MySQL leader database and then set up multiple replicas that reapply all of the changes made to that leader. It turns out the replication stream is this sort of slightly weird binary protocol, but if you know what you're doing, you can tap into that yourself and you can write your own code that reacts to changes that have been made to the database. There's a fantastic um open source Python library that we use for this at Eventbrite called um Python MySQL replication. So at Eventbrite we built a system called Dilithium. And Dilithium is essentially a way of listening to sorry about this.

27:46

Speaker 1: It's a way of listening to those um database changes and using them to trigger other actions around event the Eventbrite system. So the way it works is you have your master MySQL database with all of the rights going to it. You have a replica MySQL database that's that's replicating off of that. Then Dilithium listens to that replica. So it's replicating from a replica to figure out what changes are going on. It sees things like Event row 57 has been updated, attendee row so-and-so has had these fields changed. And it takes the um it takes that in that flow of data and turns it into more sort of um more in turns them into what we call interest. Testing moments because we can't use the word events at Eventbrite because it's already taken by one of our main domain objects.

28:32

Speaker 1: So those moments that come through are th we we translate into things like event fifty-seven has been up Updated. Event 23 has been created. Order 37 has been placed. Those we then write into Kafka, which is a very um robust uh high performance uh method queue that LinkedIn put out a few years ago. And the our search indexes are one of many different components that can then listen to that Kafka queue and decide what decide when they need to re-index things. So it's a pretty complicated flow of data once you stick it in a diagram. But um essentially what this means is anytime any piece of code Eventbrite updates one of the rows in our events table, the um Dilithium will pick that up, we'll turn that into a Kafka message. Our indexing code will listen to that, we'll say, oh, event 57 has been updated.

29:18

Speaker 1: It'll then query the database to figure out the current details of that particular. particular event and then write those changes into Elasticsearch. So the end result is we have something which scales extremely well and which can be run on many different machines at once and gives us a very robust, very robust path from initial database change to updates in our Elasticsearch Index. So I've got a few tips and tricks that I wanted to dive into. Just little little um bits and pieces that I've picked up that have helped with implementing this overall. overall pattern. And the first one is one that can really help avoid serving stale date stale data to your users. And that's to do everything with your search engine in terms of object IDs

30:05

Speaker 1: as as opposed to the raw data itself. Generally with with your search index, you'll be writing a lot of data into it. You know, it needs to know the titles of things, the descriptions of things, any fields that you might want to search. search by. And so there's a temptation to hit the search index, get that data back, and then use that to construct uh objects that you would present back to your user. The moment you do that though, you're setting yourself up some for some really nasty nasty latency risks because as we s as I said earlier, there's going to be a three to five, maybe ten second delay between changes in your database and changes in your index. And you really don't want to be showing that stale data to your users. So the trick here is very simple. When you run searches, you all you ask back from the search engine are the IDs of the underlying records. So you get you run a search, you get back a list of say 20 integer IDs, you can then hit the database directly to inflate those into actual finished objects.

30:57

Speaker 1: And it sounds like, I mean, the the downside of this is the You're adding additional load to your database. The good news is that databases are insanely quick at primary key lookups. Anytime you're doing primary key lookups or lookups against an index, that's going to return really fast. So with Django, you can use Django's in bulk help method and uh make extensive use of prefetch related as well, which again is a very fast way of retrieving data. And you can Set it up so that your users will never see stale data because that stale data, even if the data was stale in the index, by the time it's pulled from the database, it's going to be the most recent version of things. A related concept to this is what to do if there's so if you're doing this, if you're pulling things directly from your database, what do you do if something's been deleted?

31:42

Speaker 1: What do you if your search engine gives you back? ID57, and then when you hit fetch from the database, ID57 has been deleted in the time it took for you to for that that that search to come through. And the way we've handled this in the past is um essentially Essentially to have a self-repairing mechanism. The code that queries the database can notice when there's an ID that doesn't mat that that's missing, and then stick it on a salary queue or delete it or or or um or some other mechanism so that the search index knows to then remove that document. from the index entirely. The downside of this is somebody might ask for a page with 10 results on and one of those results is missing, so you you end up giving the Mac nine results and then quietly filing that tenth away to be deleted. My hunch is that no one will ever notice this. I don't think people go around counting the number of results they get on a page, so it's it's probably okay.

32:31

Speaker 1: And then one last trick, um, which again ties into this idea. This is something I've I've been I'm using on a project at work at the moment. And I'm calling it the accurate filter trick. Essentially this is a way of um solve is of uh it's an additional way of solving for this latency between your search index uh between your database and your search index index. So imagine if you will that you're building a system where users can have uh you you users can save events. So they'll see an event that they want to go to, they hit a save button and that event is is saved to their account and some way. I would like to be able to pull the to to answer queries about what events this user has saved by hitting the search index. Because if I can hit the search index, I can combine it with all of these other benefits.

33:17

Speaker 1: I can let users search with it, search for text within the events they've saved. I can do filters by geography. There's all sorts of useful things things I can do with this. And that's easy enough to implement. You have a field on your event document in Elasticsearch with containing the IDs of the users who have saved that event. And then you can do a search like this. You can say search for events where one of the saved by users values is the user ID that I'm dealing with. There's just one obvious problem with this. If a user saves an event and then goes and looks at their list of saved events within a few seconds and it's not the and and that event isn't shown to them, then obviously something is broken. You know, that that they this is the latency problem that we've been fighting since since the beginning of the talk. So what you can do is you can say, okay, anytime I'm running that query, the first thing I'm going to do is hit my relational database to figure out what are the events that

34:08

Speaker 1: This user has saved in the last X minutes, let's say the last five minutes. So this is guaranteed to give me an accurate model of the user's recent activity, and it'll give me back, say, four or five um IDs of documents that we know that That the user has saved. Once you've got that list, these are the ones that were saved in the last five minutes. You can construct an elastic search query where you say, give me back any event where either the user is listed in the list of users who have saved this event Or the event itself was is one of these five that we know that they've saved recently. And as a search query, this will run crazy fast. It gives you all of those benefits, but this is guaranteed to be exactly up to date with the activity of your users. Saved by users is one obvious um application of this.

34:53

Speaker 1: There are a whole bunch of other things where if you want to get a precise up-to-the-date reflection of the state of your system you can use tricks like this to pull those out of Elasticsearch. So a few more use cases that I've um applied in in particular I've applied Elasticsearch to. One that we use at Eventbrite is for recommendations, because it turns out recommending events to a user is essentially just another search problem. You can there are a bunch of ways you can do this. You can say find events where one of my friends has saved that event. This is the way we did our calendar for Lanyard earlier. You can also say find events that are similar to the last ten events that I've saved. A very straightforward way of doing that is to look at the last 10 events saved by a user, collect together the text from the title and description of all of those events.

35:42

Speaker 1: events into a giant blob of words, and then just search for those words with a Boolean or clause, which is enough for Elasticsearch 's relevance to kick in, and it'll give you back other events that are that are similar in textual content to the events that the user has searched. saved as well. And search engines are really good at relevant scoring and boosting, so you can fine-tune this stuff very to a huge, huge extent using the tools that are built into the search engine. engine. Another thing Elasticsearch is great at is geographic search. It's got built-in support for Geo. You can add latitude and longitude points to your documents and then you can do things like all event or all all documents within five kilometers of this radius point. You can even send it a polygon and say, here is the shape of Canada, give me everything that falls into

36:28

Speaker 1: That polygon shape. And again, this is stuff which, if you're not using Postgres , can be quite difficult to do with a relational database. But more importantly, you can combine these with all of the other search and filters. So if I've got a recommendation system built on Elasticsearch, I can say recommend me events similar to these events that fall within this geographic area and further combine that with other options as well. Elasticsearches and search engines in general are great for visualizations. You saw this earlier with the DC Inbox Explorer. I've got this little graph at the top, which is actually generated from yet from just another one of these aggregations. When I search against Elasticsearch, I can say give me back counts per month for this time period. And once I've got those counts back, I can turn them Those into a bar chart.

37:13

Speaker 1: And if you look at the way Elasticsearch is used for log analysis, people do some really exciting visualizations on top of the raw data that's being collected by these aggregations. And one way to think about this is it's kind of like having a real-time MapReduce engine. If you've ever used MapReduce on something like Hadoop, it's a very powerful way of running a query across many different different machines and getting results, but it's generally something you want to run as a batch job because it might take 30 seconds to a minute for it to return results. Elasticsearch under the hood is doing pretty much exactly that. It's got if you give it a search, it will spread it out across the c the nodes in your cluster, combine the results together and use that to return return um return documents to you. But it's designed to work in real time. You're getting like

37:58

Speaker 1: um response times measured within milliseconds, which means that you can expose these directly to your users. So in summary, you should denormalize your data to a query engine. It's definitely a good idea. It lets you build all sorts sorts of things you couldn't build before. And Elasticsearch, it turns out, is a pretty good option for this. And I've left lots of time for questions, so thank you very much.

38:22

Speaker 2: Hi.

38:23

Speaker 1: I will take that question.

38:26

Speaker 2: Thank you. You were talking about how you could use a queue to sort of de-duplicate repeated index reactions. Can you speak a little more to that?

38:34

Speaker 1: Sure. So the thing that you want to avoid is 500 people like 500 people interact with something in your database and then you send 500 updates to your Elasticsearch index in a in a giant flurry. So really this is um is it de duping? Yeah it's deduping it's rate limiting um it's a It's being able to get smart about this and say, okay, there were five hundred updates within a short space of time, but I'm actually going to turn that into a single combined update to the end. index. The way I've um built this on top of a queue is your the the code that listens to the queue. Well so actually one way that I built this um against Redis was to say when an indexing request comes in If the indexer hasn't done anything in a few in five seconds, just index that thing straight away.

39:21

Speaker 1: If the indexer has run within the past five seconds, that suggests that there's a lot of activity going on. So then hold out for a couple of seconds to see if other updates come in for the same event ID. And if they do, bundle those together and send that all at once. I think actually the um the de duping becomes a lot easier if you use the last modified timestamp, because then it's just your cron job runs runs once a minute, and if there were 500 updates to an event in the last minute, you'll still only re-index it once.

39:47

Speaker 3: Hi. Um I was wondering, are you the only person who's like put this into terms and uh is Educating people about this?

40:00

Speaker 1: As far as I know, I am, which I find really surprising because I've seen lots of places doing exactly this. I just don't think anyone's put a name on it. on it before. So if somebody does have an alternative name, then I'd love to hear about it. But yeah, as as far as I can tell, I'm the only person who's said let's let's give this a name and and start discussing this as a as a general strategy.

40:19

Speaker 3: So there's no bugs.

40:20

Speaker 1: No buts or bugs?

40:22

Speaker 3: Bugs.

40:22

Speaker 1: Bugs. Oh books. Definitely not yet, no. I I should write a blog entry.

40:29

Speaker 4: Yes, please write a blog entry. That would would be great.

40:32

Speaker 5: Um so you talked a lot about uh getting stuff out. I was wondering uh like specifically for your DC mailbox example, do you have to do a lot of pre-processing of the data before it goes in when you're actually structuring your documents and stuff like that. Is there an extra step on that side?

40:51

Speaker 1: Yeah, so um one concept I didn't really talk about is um there's a thing in Elasticsearch called a mapping, which is basically the same thing as a SQL scheme You don't have to use mappings. You can just start blasting JSON documents at it and it'll it'll work. But if you actually want to be able to differentiate date and times from geographic points and so forth, you need to use that. And so here's the mapping for the DC inbox thing. And actually this one's there's a lot of data here, so it's got caucuses and Congress numbers and C span IDs. And so forth. The source code for this is all available on GitHub, so you can see. be a slightly higher level way of working with Elasticsearch.

41:37

Speaker 1: You can also compose this as a JSON blob and then post that to Elasticsearch itself. But yeah, so this is your first step is going to be design a mapping for your data. The actual indexing is pretty trivial once you've got the mapping in place, because you really are just constructing JSON documents and then posting them back up. of the server. But yeah, the mapping design is is quite important.

41:56

Speaker 6: Thanks for the talk. That was really great. So I I know that you're talking about how performant this is to find documents based on whatever criteria. So is it equally performant if I wanted to do something like aggregate so to get like average scores or standard standard deviations. on something inside the document?

42:12

Speaker 1: Absolutely. This is the the strength of aggregations is that they're insanely fast for stuff. I mean I it's As you get more complicated with your aggregations, the performance can it can start to add up. But honestly, the most complex queries I've come up with are in the order of 100 milliseconds, I suppose, to 10 milliseconds. So generally the performance is really good. And yeah, there are a lot of the aggregations I've shown so far are what are called bucket aggregations where you divide documents into different named buckets. There are also um aggregations for metrics like that can capsulate things like standard deviations and sums and medians and uh even like in geospatial terms there are aggregates that will calculate a bounding box around all of the documents. So there's a whole bunch of additional power and flexibility I suppose you get around that.

42:56

Speaker 4: All right. Thank you so much, Simon.

42:58

Speaker 1: Thank you.

Questions this talk answers

What is the denormalized query engine design pattern?

Keep the relational database as the system of record, copy selected data into a separate search index, and synchronize changes between them. This combines reliable relational storage with the search engine’s strengths in counting, aggregation, complex queries, and horizontal scaling.

Discussed at 1:46

Why use a search index alongside a relational database?

Relational databases become expensive for large counts, queries that scan many rows, and complex conditions spanning multiple fields. Search engines are better suited to fast counting, aggregations, multi-field queries, relevance scoring, and full-text search.

Discussed at 3:18

How did Flickr use a search index to solve queries across sharded databases?

Flickr stored users’ data across database shards, then loaded all photos into a horizontally scalable search index so queries such as finding every photo tagged “raccoons” did not need to hit every shard.

Discussed at 8:44

How do you avoid showing stale search results after a user edits their own data?

Route requests for a user’s own data to the relational database, while using the search index for public or other users’ data. This avoids exposing the search index’s synchronization delay where the user can immediately notice it.

Discussed at 10:17

How can Elasticsearch provide faceted search and nested aggregations?

Add aggregation definitions to the JSON search request to count matching documents by fields such as party or role. Aggregations can be nested, such as counting role types within each party, and can support highly interactive filter-and-drill-down interfaces.

Discussed at 21:09

What are effective ways to synchronize a relational database with a search index?

Options include polling indexed rows with a last-changed timestamp, sending changed IDs through a queue, or consuming the database replication log and publishing changes to a durable system such as Kafka. The approaches become progressively more sophisticated and can support batching, replay, and scaling.

Discussed at 23:00

How can you prevent stale data when using search results in an application?

Have the search engine return only the matching object IDs, then fetch the current objects from the relational database using primary-key lookups. The index may be stale, but the database fetch supplies the latest data.

Discussed at 30:05

What should happen if a search result refers to a deleted database record?

Treat the missing record as a signal to repair the index: omit it from the current response and enqueue it for deletion from the search index. This lets the application self-heal from index/database inconsistencies.

Discussed at 31:42

How can you make searches involving a user’s recent changes accurate despite index lag?

Read the user’s changes from the relational database for a recent window, such as the last five minutes, and include those IDs in an Elasticsearch query alongside the indexed criteria. This preserves search and filtering benefits while guaranteeing that recent activity appears immediately.

Discussed at 33:08

How can Elasticsearch be used to build event recommendations?

Recommendations can be expressed as searches, such as finding events saved by a user’s friends or finding events similar to the user’s recent saves. Combining text from saved events and using Elasticsearch relevance scoring can return similar results.

Discussed at 34:53

Presenters

Note: We understand that names change, people change, and bodies change. We respect each individual's journey and privacy. If you have any concerns about a video or need us to remove content, please don't hesitate to contact us. We will handle your request with care and promptly address any issues.

More videos by Simon Willison

More videos from DjangoCon US