Bulletproof Data Pipelines: Django, Celery, and the Power of Idempotency

This video features Ricardo Morato Rocha at DjangoCon Europe 2025 in Dublin, Ireland.

Bulletproof Data Pipelines: Django, Celery, and the Power of Idempotency
0:26:16
Published June 4, 2025
1,029 views

Talk: Bulletproof Data Pipelines: Django, Celery, and the Power of Idempotency by Ricardo Morato Rocha

https://pretalx.evolutio.pt/djangocon-europe-2025/talk/BRWVBM/

Summary

Ricardo Morato Rocha presents a step-by-step redesign of a video enhancement pipeline that processes every frame through an external API. Starting with synchronous, failure-prone code, he adds exception handling, Celery tasks for parallelism and retries, and Django models and managers to track each frame’s state. He argues that idempotency—implemented here with a status check and database locking—prevents duplicate work, reduces API costs, and makes failed frames independently retryable. The broader lesson is to evolve architecture incrementally, keep workflow state in the database, and treat Celery primarily as a message queue rather than the owner of complex workflows.

Key takeaways

  • Adding per-frame exception handling prevents one failure from aborting an entire video pipeline.
  • Celery can parallelize frame processing and support retries, but complex workflow state is often easier to control in Django models.
  • A database-backed execution record shows which frames succeeded or failed, allowing only failed work to be retried.
  • Idempotency checks and `select_for_update` prevent duplicate API calls and race conditions among workers.
  • Resilient architecture should evolve in small, practical steps while acknowledging trade-offs such as database volume and third-party API load.

Summarised automatically from the transcript.

Chapters

  1. 0:00 Introduction and Thought Experiment Ricardo introduces the talk’s focus on resilient background processing, incremental refactoring, and software architecture.
  2. 2:26 Django Vids Case Study The example application and its growing video-enhancement workload are introduced.
  3. 4:46 Video Processing Pipeline The talk breaks down the multi-step upload pipeline and identifies video enhancement as its bottleneck.
  4. 5:35 Synchronous Frame Enhancement The initial frame-by-frame implementation is examined, including its failure behavior and scalability problems.
  5. 7:14 Incremental Error Handling A small refactor adds exception handling so one failed frame does not stop the entire video.
  6. 8:45 Celery Task Parallelization Celery is introduced as a distributed task queue for retries, asynchronous processing, and parallel frame enhancement.
  7. 10:20 Execution Observability The limitations of Celery’s visibility are discussed, motivating database-backed execution records and status tracking.
  8. 11:07 Idempotent Frame Processing Idempotency is explained as a way to skip frames that have already been enhanced and avoid unnecessary API calls.
  9. 14:13 Concurrency Control The execution model is refined with database status checks and row locking to prevent duplicate work among parallel workers.
  10. 18:22 Extending the Pipeline Pattern The same execution strategy is applied to other pipeline stages using a chain-of-responsibility approach.
  11. 20:26 Scaling Trade-offs and Lessons The talk concludes with practical trade-offs, API-load concerns, incremental architecture improvements, and the importance of idempotency.

Transcript

4,172 words · auto-generated Show

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

0:06

Speaker 1: Hey, everyone. I'm Ricardo and it's okay. There are many different pronunciations for this word. There's like idempotency, identity, a lot of different pronunciations. Before we start, just a quick question. Were by any chance any of you here less on last year's DjangoCon? Uh because there was a talk by Jake Howard, if I'm not wrong. uh talking about salary and workers, background workers. And it was a was a terrific talk, very good. uh I truly recommend you watch if you didn't uh it's available on a youtube and I'm mentioning this because I believe that this talk is like a spiritual follow-up from this

0:51

Speaker 1: because we will be talking about salary, but could be any kind of messaging queue system actually, and how to make it resilient enough, but still managing complexity so it doesn't explode in our faces. Uh but okay, before we start, as I said, as we said, I'm Ricardo. I'm born and raised in Recife, Brazil, it's the northeast part of Brazil. I spent some time in Portugal, uh here in Europe, in Portugal. I'm a software engineer at Vinta Software, which is a company focusing focused on projects and software projects, uh basically fully focused on Python and everything Django related Uh and I'm constantly talking about architecture, both real-world and yes, software architecture.

1:37

Speaker 1: Maybe because my fiance, she's an architect, like a real-world architect that builds things, so I just like to talk about it. And this whole talk is just an excuse for me to talk about architecture for like 30 minutes. But okay, what is this? What is this talk? actually this talk is a thought experiment it's not it's based on a real world case scenario that happened to me and my team so we'll be following some strategies that we use and explaining how how we do this But most importantly, it is a just a thought experiment for better understanding parallelization and software architecture and also how we can refactor things in a organic way. So we you start with a piece of code. that is just bad. I think it's just like a method that breaks up every so often and and

2:26

Speaker 1: isn't very resilient and we will iterate through it. making the bit better just bit by bit until in the end of the talk we have something somewhat okay but with with still room for improvement And welcome to Django Viz. Django Viz is a video platform aimed for the Django community where everyone can upload their uh tutorials, maybe their videos. Uh it's a new platform. With hundreds of users, new ones coming every day, and thousands of videos. It is a very strong community. Uh and new uploads are frequent in the in the platform. They are occurring every minute. And for each video, we we run a enhancing pipeline

3:12

Speaker 1: where the video will be like the the every frame of the video will have better sharpness, better contrast, better brightness, something like that. And the Jungle Vid's uh team is saying that scaling is hard. The pipeline, this video enhancement pipeline is getting out of hand, it's taking too much time to actually do the things. And scaling is hard and that's normal, right? I mean we all have problems with scaling, we have a lot of different talks about how we can m better manage, I don't know, database, how we can better manage our software, which tools we can do to scale better. But is it really hard? Here is a Google Trends picture, and by no means this is like a scientific researcher

3:58

Speaker 1: or something like this. But this picture shows that searches for architecture, software architecture, are decreasing in Google since 2004. Maybe because we already have like resources to handle those scenarios and those kinds of problems. Also have several books and articles published in this topic. thousands of large-scale distributed systems are operating worldwide. I mean Amazon, Netflix, Facebook, Uber, they all handle this constantly in every single second. They have a lot of different pipelines. Everything's running smoothly for them. Why not for JungleVid? That's what we'll be looking at right now. Let's take a look at the video upload pipeline. And when I say pipeline, I apologize in advance for any data scientist or data engineer.

4:46

Speaker 1: I'm I'm just simply talking about a pipeline as a F is a multi-step process. So as we can see, once we upload the video, the raw video data goes through the first step, which is pre -processing, where we get the metadata, uh metadata, the we decode something. I don't know, something like that. Then goes through the video enhancement process, which is the most important and is the bottleneck of this whole operation Then it goes through a thumbnail generation step and then a speech to text caption creation. And it outputs enhanced video that will be available on the on the platform. And as I said, this video enhancement is the bottleneck right now. And why is this? This bomb code uh is a lot of code, but it's a simple code.

5:35

Speaker 1: And this is the thing that I said, it's just some somewhat bad code. So for each frame in the video, we're making an API request. We are checking for the the status quo. If the status code isn't okay, we just explode everything, raising an exception. But if if it is okay, we enhance the frame, save the frame, and go to the next one. I mean there are a lot of problems here, just a lot of problems, but do not uh go through them right now. We just take a bit uh uh just a step to make a bit better But beforeward, this this visualization may help to understand. So we start the enhancement, we iterate through each frame of the video And for example, if we have a one-hour long video with 30 FPS frames per second, it's more than a hundred thousand frames that will be iterating through for just one video.

6:27

Speaker 1: We make one API call for each frame, we check this API response, we actually don't not don't handle the failure, just explode everything, and then we save the frame if if it uh if it's okay. And this exception handling is a volcano about to erupt. So this is the first pain point that we can see in this in this whole process. Uh we do not have any type of error handling. Uh we have a fully synchronous API, so we go through each frame one by one, taking a lot of processing time, because remember we do have to make an API request for each frame. So let's make it better just a bit. And as I said, this is going to be like an organic refactor, more like a small step to a better outcome.

7:14

Speaker 1: And the only change that we'll be making is putting a tracks at block. This isn't isn't like some magic or some extra fancy thing, but it does handles one of our pain points. So now we can ensure that at least we'll try to enhance every frame. Before if frame number two uh failed for some reason, the whole thing will would explode and we We're gonna go through frame number three, four, and so on. Now actually at least we try to go thro through them. Uh but it doesn't still uh it doesn't still have a lot of pain points. How can we retry those failed enhancements How can we paralyze those enhancements and make them f uh uh I don't know maybe faster? Uh and

7:59

Speaker 1: Here is uh another visualization of this. So we saw the enhancement iterate through the frames, attempt now to f the to do the enhancement of the frame. Uh we actually handle the enhancement error now and then continue to to neck the next frame As I said, it isn't a lot of refactoring right now, but it it is just a bit better. And this is the whole idea of what we're doing here is to keep evolving bit by bit, PR after PR, to make just a little bit better. How can we keep improving this? We can handle API and availability, we can retry field enhancements, we retry some some frame that for I don't know maybe the API was not available for uh like two hours

8:45

Speaker 1: and we had a lot of uh uh frames failing to enhance how can we try those? How can we parallelize those enhancements And for this specific pain point, we can use salary. So salary and it could be like our queue, it could be actually any message system Will help us make it a sync, which is another small step for our maybe more resilient implementation Celery, for those who don't know, it's a very popular package in the Python community and in the Jugo community as well. It's a distributed test queue, it's highly available, it's horizontally scalable, which is which will be very important for us. And it is real it is a reliable distributed system to process vast

9:31

Speaker 1: amounts of messages, just like our use of case, because as I said, if you have an one hour long video, we'll be handling more than a thousand messages. for this video alone. So this is the next step we can do to improve the enhancing algorithm. Now instead of actually Going through each frame and trying to enhance the task. Now just enhancing the frame, now we just call a task and queue a lot of different tasks, one for each frame, to try to enhance it Uh salary already has some an availability and auto-retrice configurations. As per my experience and what we had on my team. It is it isn't so intuitive to know to to do things on salary, those kinds of

10:20

Speaker 1: things in salaries, so it is still a bit complex. And also how can we pinpoint why the frames failed? How can we improve this and have better uh observability without the need of an extra tool, for example, without the need of data dog, new relic or something like this, how can we with a jungle query just get get this from our database? This is uh a balance between the benefits and challenges from salary. So it's very good that now we have an async process and we can have the auto-retry capability But it we don't have a lot of visibility.

11:07

Speaker 1: That's where the manager is in id important executions or id important executions enters the room. Now we'll unlock the true power using Django's ORM. So it'll be as simple as a module. uh is doing just a query on a terminal to get the data from this model and why it's important executions. Basically from this extra uh excerpt from Wikipedia uh neat important execution is anything that can be applied multiple times without changing the result be beyond the initial application. So imagine we have a frame that we want to enhance. If this frame, uh if we enhance this frame once, we have like better sharpness, better uh noise reduction, and etc.

11:52

Speaker 1: If we enhance this frame again, we may have even better sharpness, even better noise reduction. Uh I don't know if you have any familiarity with uh with this kind of API with enhancement photo enhancement APIs but basically they always find something to make a bit better And it so it isn't natively indepotent. If we keep making the API requests, if we keep uh reaching out to the API, it will keep enhancing and enhancing and enhancing, and some sometimes it's not even You can even perceive a difference, but there is a difference. So how can we ensure we ensure that we'll not be trying to enhance frames that were already enhanced? uh skipping unnecessary API requests that will save us time and also money when we don't reach out to this

12:41

Speaker 1: uh expensive API. Here we entered the executions realm. Maybe uh basically you have a frame enhancement execution, which is just a jungle model. It's just a jungle model that has a manager. We'll talk about it later on. It has a frame, the frame from the video you want to enhance and has a status. This status will be very important for us to have granularity over this the status of each frame. So we'll know which frame failed, which frame is running, and so it's been under the process of enhancement. which frames were not enhanced, uh and so on. So it'd be very important once you have like a parallelized uh tasks running in parallel with

13:27

Speaker 1: a lot of distributed workers. Uh and it would be very easy to pinpoint which frame are failing, which frame are are uh painting to be enhanced Taking a look now at the manager, the V enhancement manager will have a list of frames, which uh we saw in the last slide uh it it all it is as simple as a jungle model can be it's just a model that has a foreign kit video So we can keep the status the track the status of the video enhancement and also has a status so we can we will know if the manager has started, if the manager is running, uh how many frames of of the video are no are being uh enhanced right now and etc. And here

14:13

Speaker 1: is Our updated algorithm. So basically we create the manager, we create all the executions for for this manager based on the video frames, and we once again delay the task. So that's not much different from what we are doing right now. But now with this kind of implementation, we can have, for example, a sidecar. or in the in in the background running uh just for the failed uh frames we can keep retrying those frames every five seconds, every ten seconds, every minute. uh we have a lot of room for improvement here that we didn't have previously. So if we take a look b uh a step back and take a look at what we've done here, we've came from a simple method

15:01

Speaker 1: that was fully synchronous, it broke every so often uh and wasn't reliable. Now we still have a method that it can break. Uh it may not be fully resilient, it's fully reliable, but now we can understand which frames were broken, which frames uh we couldn't uh we couldn't enhance for I don't know. Some reason, uh maybe API in availability or something like this. Uh we can retry those frames without the need of running the for loop again for the whole video, trying to enhance again frames that were already enhanced So when we we think that each enhancement here is a API uh call that we do, every s every frame that we skip on our for

15:47

Speaker 1: loop is one less API call that we're doing. And zooming in right now on the task itself, we basically just select the execution, get the frame from the execution. We still uh now we have just a status check, and this is where the independency enters. As I said before, enhancing a frame isn't natively indepotent. We will always try to enhance even if it has been already enhanced. So it's just a status check. Again, it's not it's not something fancy, it's not like a finite state machine or something like this. It's just a status task. uh status check if the enhancement was already completed we just skip this frame there's no point of enhancing it again

16:33

Speaker 1: so We do this this check, try to enhance the frame. If the video enhancement uh throws any type of exception, we get this exception, we save the execution as failed. And we move on. We just go through to the next frame that is running in parallel. So now we have a lot of workers in parallel. We can have failed frames as we will, uh because We don't need to like retry every frame, just the failed ones. We can have a sidecar, we can have another uh design pattern to you to just get get those frames. From a simple database query, you can just square the frames that will fail for some reason and retry those frames. So It's again, it's just a bit better than was before, and there's always room for improvement.

17:22

Speaker 1: There's always something better that we can do. Uh also here in the in the top where we get the fra the execution, the frame enhancement execution, we do a select for update because we are now working parallel and sometimes when you have a lot of different Uh workers, you may have a worker that is speaking up a frame A and uh is another worker is speaking up the same frame So this will will create a race condition. We don't want this. It's it's a concern that we it's a more complex approach than before, but it's also more resilient approach. So this is why we need this select for update line uh and this is now our algorithm our process funnel

18:08

Speaker 1: So we create the manager, the manager has several executions, one for each frame of the video, we delay several tasks, one for each frame of the video, and we keep checking those status. We can run the enhancement process for the video ones. Maybe it will fail. Then we can run the pro the process again and just pick those that have failed and do not try to enhance those frames that were already enhanced. And if everything goes as expected, we'd have a enhanced enhanced frame at the end. And A question is can we keep improving? This is as I said, this is the main mantra, the the the the whole point of this whole talk is There's always room for improvement.

18:54

Speaker 1: There's always something we can do to be a bit better, some step you can make to go in a different direction. So we could save our metadata in the failed executions, for example. We could know why this enhancement failed. Maybe just a JSON field in the database would be enough. Postgres has a JSON field for this. MySQL has as well. So maybe this is most that we can do to understand why uh what we can make better for the API requests, for example. What we can make better what we can do to be better to be more resilient Also, we can implement it in other parts of the pipeline. As we saw in the in the beginning, Django Vids has a multi-step pipeline. This is only one piece of it.

19:41

Speaker 1: So how can we use this for other steps? How can we use the same strategy for other parts of the pipeline? Can we use it? In my team, we do use it. It's something that is working for us, especially when you use the chain of responsibility, where it's a principle that we have A step that starts and runs and ends, and then when this step ends, we start the next step and and so on. So the whole pipeline could be just a big chain of responsibility where in the first step the metadata we run us uh a lot of different workers to collect different metadata types from the the video Then when all these workers are successful, we go to the next one, which is a video enhancement. This one that

20:26

Speaker 1: we just saw. Then you go to the next one, to the next one, to the next one. And it's actually just something Really simple to implement when you get grasp of it. And again, it's just as simple as a jungle model can be. And finally, what can you get from this experiment? Better architecture design may actually will cause degraded performance. We haven't changed a lot in the architecture of the API, the a lot of the architecture of the code that we are dealing with We just basically did some exception handling and made it a sync. But yeah, that's step by step, bit by bit, just making a bit better. Uh we j we could, for example, no not make an API request for each frame, because we are now having a hundred thousand different requests for each video.

21:18

Speaker 1: So yeah, this could be way better than all this uh that we did, but it could also be a big refactor that the team's not ready to do. So it's always a trade-off to understand what you can what you cannot do right now. Uh properly scaling a simple pipeline as this should not be a pain. We just saw it's like one step, then another, then another, very simple steps. It should not uh take 30 minutes to run, it should not take like any even fifteen minutes to run. Uh and with salary and with a sync workers, we can just run more workers in parallel. We can have 10, 15, 30, 100 workers running parallel, picking up frames are running, and etc. The architecture should always be evolving as more challenging scenarios appears.

22:05

Speaker 1: So we got the problem from the beginning. We keep evolving and evolving and evolving until we have somewhat something resilient, but maybe not. Uh it may be not optimal yet, definitely it isn't optimal yet, but we can still improve it and we can keep improving it as more challenging scenarios appear. For example, now we have a hundred workers working in parallel How does the API handles 100 requests coming parallel every second, every millisecond? It's a new challenge that we need to face it as a team. Uh sync workers are best friends when dealing with large amount of amount of data, but they can also be our worst enemies when dealing with a large amount of amount of data, making, for example, an external API, a third-party API.

22:52

Speaker 1: having a degraded performance because now as we said I have a thousand workers reaching for this API every second. We are just doing like a DDOS on this on this third-party API And it impotency is key when dealing with these scenarios, even when artificially created, as we did here, because as you said, the enhancement is isn't natively uh either potent but with a simple status check we can make it somewhat either potent and make sure that we'll not be spending CPU cycles and API requests and and etc in something that cannot be enhanced even further. And that's it. Thank you.

23:41

Speaker 2: Thank you, Ricardo, for the talk. So uh my question is. When you started doing the uh optimization, uh were you trying to save each frame details or frame into uh each row of the table. So like if you are having one GB of uh video uh and you said probably we are considering 30 frames per second So were you having that many frames being recorded in the uh sequential database?

24:16

Speaker 1: It's a great question. For this example alone, yeah, it happens like this, but this is As the first point of the final thought said just bad architecture. It should be just like this. It doesn't make sense to do something like this. Yeah, it's something that we could improve for sure. Uh maybe it's just a U refactor that is way too big. for f to do for the small team, for example, and doing small steps to improve the the the pipeline could do a better could be a better approach, maybe. But yeah, I completely agree. It's just uh not a good architecture to begin with.

24:57

Speaker 3: Thank you for the talk. From your experience, which is the biggest salary pain point from your experience? Something that grinds your gears?

25:07

Speaker 1: That's a great question. And I will say it is workflows. having complex workflows with salary. Uh it was the actual the the the seed that started all these movements from us to move away from hand uh having salary handling everything like workflows and chains and etc and having something more jungle focused approach so Complex workflows with one one uh task depending on the other. When uh task number three fails, what happens to task number four, five, six? It's something quite complex to understand, uh at least in previous versions of salary that we were working with. So moving away from

25:49

Speaker 4: charts and groups, where in controlling the status in the database, that's what we're affecting.

25:54

Speaker 1: Exactly, exactly. Moving away from anything like complex workflows on on salary, using salary just like a message queue, and that's it. But controlling the state on database was our goal. Thank you, Ricardo. Thank you everyone.

Questions this talk answers

How can I make a video frame-processing pipeline more resilient with Django and Celery?

Refactor incrementally: catch errors so one failed frame does not stop the job, queue frame enhancements asynchronously with Celery, and track each frame’s execution status in Django so failures can be retried independently.

Discussed at 7:14

How can Django models improve observability for Celery pipelines?

Use Django models for the pipeline manager and per-frame executions, including statuses such as pending, running, failed, and completed. This lets you query which frames failed or are still running directly through the ORM, without requiring a separate monitoring tool.

Discussed at 11:07

How do I retry only failed Celery tasks without rerunning successful work?

Create a database execution record for each frame, mark failures explicitly, and have a later process or sidecar query and retry only failed executions. Completed frames are skipped, avoiding another full pass through the video.

Discussed at 14:13

How do I make a non-idempotent frame-enhancement API effectively idempotent?

Store the enhancement status for each frame and check it before making the API request. If the frame is already complete, skip it so repeated pipeline runs do not spend more time or money enhancing it again.

Discussed at 15:47

How do I prevent multiple Celery workers from processing the same database record?

Use `select_for_update` when selecting the frame execution, so concurrent workers cannot claim the same frame and create a race condition.

Discussed at 17:22

How can I apply this reliable-task pattern to a multi-step video pipeline?

Treat the pipeline as a chain of responsibility: run parallel workers for one stage, start the next stage only after the current one succeeds, and use the same execution and status-tracking strategy for each step.

Discussed at 19:41

Is storing one database row for every video frame a good architecture?

No. The speaker agrees that recording every frame this way is poor architecture, but says a large refactor may not be practical for a small team, so incremental improvements can be a better short-term approach.

Discussed at 24:16

What is the biggest pain point of using Celery for complex workflows?

Complex task dependencies and workflows are difficult to reason about, especially when a task fails and later tasks depend on it. The speaker’s team chose to use Celery primarily as a message queue and control workflow state in the database instead of relying on Celery’s chains and groups.

Discussed at 25:07

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 from DjangoCon Europe