Пет Проект для Data Engineer | Инженер Данных — Transcript
Full transcript
- 0:00Alright, recording is on. Hello
- 0:01everyone. Anyway, today we have a
- 0:02project for a data engineer. Where did
- 0:04it even come from? We have a bootcamp.
- 0:06Obviously, we launch new cohorts every
- 0:082 months. And in the current cohort, we
- 0:12have a really cool guy. His name is
- 0:14Artyom, if you please. He is currently
- 0:16available here via video link. He is
- 0:18real, not an AI or anything like that.
- 0:19We really have live people here. He did
- 0:21, well, you know, we came up with tasks
- 0:24, and people usually do the tasks as
- 0:27written. But Artyom went further. He
- 0:30really wanted to figure things out, so
- 0:31he did a pet project. Meaning a
- 0:34pipeline that's broader and bigger,
- 0:36both in terms of tools and logic, and
- 0:38processing. And it turned out to be
- 0:42really rich. We have many tools:
- 0:44Greenplum, ClickHouse, dbt, Trino, and
- 0:46so on. But Artyom actually put together
- 0:49such a cool project that's worth
- 0:51checking out. You know, by the way,
- 0:54what I say about data engineers? We
- 0:56have these tool stacks. I've started
- 0:59using an interesting phrase lately,
- 1:02like "a cool stack." It's like, you
- 1:05know, P2P arbitrage. Something like
- 1:07that. Only for us, it's a stack of
- 1:09tools. And today Artyom will tell us
- 1:11how he connected them, where he gets
- 1:13the data, who the client is, who the
- 1:15consumer is, what ETL is, data layers
- 1:17—those ODSs, stages, and so on. dbt,
- 1:20what is it used for? Greenplum,
- 1:22ClickHouse, Spark, Airflow—in short,
- 1:24the full kit. So, we're all watching
- 1:27today. I'm interested, too. Artyom, the
- 1:29floor is yours. Tell us a bit about
- 1:32yourself, like who you are, what you
- 1:34stand for, and so on.
- 1:36Yeah, Zhenya, hi. Volodya, hi. Hi
- 1:39everyone who joined. So, my name is
- 1:41Artyom. For the last year and a half,
- 1:44I've worked in a position very close to
- 1:47data engineering. My company was
- 1:49involved in implementing data analytics
- 1:52and BI systems, but I wanted to become
- 1:55a pure DE, so I joined the guys at the
- 1:58bootcamp. I already knew some
- 2:00technologies, of course, but I wasn't
- 2:02familiar with others at all. And plus,
- 2:06I even improved what I already knew
- 2:08thanks to the guys. So today I will
- 2:11show you my project, essentially as the
- 2:14finale of what I've learned here. I
- 2:17think it will be interesting. And feel
- 2:19free to ask questions today. We'll look
- 2:21at the presentation, and hopefully,
- 2:23we'll scroll through the code at the
- 2:24end. Yes, go ahead and show the
- 2:26presentation, tell us about it
- 2:27theoretically, and we will answer
- 2:29questions in the chat if anything comes
- 2:31up. If there’s something really
- 2:33relevant or urgent, I’ll ask you,
- 2:36just in case. There. But if it’s
- 2:38something really cool, something
- 2:39special. Otherwise, guys, if you have
- 2:41any general questions, let’s save
- 2:43them for the end.
- 2:43Well, so that’s how the project looks
- 2:46overall. And the project is indeed big,
- 2:48with a lot of technologies. It’s
- 2:51important to say here that the project
- 2:53is more of an overview. Basically,
- 2:55during the bootcamp, I learned a
- 2:57certain stack, understood how to link
- 2:59things together, and managed to write
- 3:01this. In general, the whole project can
- 3:04and should be organized as separate
- 3:06DAGs, which would look more like a
- 3:08production scenario than one big system
- 3:11. I mean, don't think that at work
- 3:13you'll be asked to do something this
- 3:14grandiose. Most likely, a quarter of
- 3:17this. Well, there are benefits to the
- 3:19project being this big. We’ll
- 3:20practice several real-world patterns
- 3:22right here. For example, how to
- 3:25transfer data from S3 to Greenplum
- 3:27using Spark. And how to transfer data
- 3:29from Greenplum to ClickHouse, also
- 3:31using Spark. And how to build models in
- 3:34DBT? And how to load it all into
- 3:36Airflow later? There are many such
- 3:39real-world tasks gathered in one
- 3:41project. I think it will be very
- 3:42interesting. Let’s go through the
- 3:45whole project in general, and then
- 3:47we’ll talk about each component. So,
- 3:49I first broke it down into key
- 3:51components, implemented them separately
- 3:53, and then combined them. You can see
- 3:55the arrows on the slide. Overall, each
- 3:58arrow is a separate task, you could
- 4:01even say a separate DAG. I wrote many
- 4:04DAGs first, and then connected them
- 4:06into one. I wouldn't have been able to
- 4:08write it all at once. Airflow is a task
- 4:11orchestrator. It can run tasks on a
- 4:14schedule. Well, it’s generally
- 4:16described here what Airflow does.
- 4:18I’ll just add that this isn't all we
- 4:21need to get from Airflow in this task.
- 4:25Take a look, we access an API. And it
- 4:29requires a token. You can’t just put
- 4:32the token in the code; it’s not safe.
- 4:35Airflow has its own storage for such
- 4:37values, called Variables. We store them
- 4:39in Airflow as key-value pairs, so we
- 4:42access it by key and get our token. If
- 4:44the code leaks, no one will find out
- 4:46the token. We can check that in the
- 4:48code later. The same goes for database
- 4:50connections. ClickHouse, Greenplum, and
- 4:53alerting in Postgres. All databases
- 4:56require a host, port, login, and
- 4:58password. You can’t store all that in
- 5:00plain text either. For such things,
- 5:03there is also another mechanism, called
- 5:05Connections. It’s very similar to
- 5:07Variables, but it returns a connection
- 5:09object instead of just a value.
- 5:11And what is that? Where is the data
- 5:13from? What kind of data is it? Where do
- 5:14we store it? At a very high level.
- 5:16Let's just run through it.
- 5:18Our data comes from newyorktimes.com;
- 5:20it’s a site that aggregates news and
- 5:22writes its own articles. We take the
- 5:24top most popular articles for the week.
- 5:28And sometimes this data changes, so we
- 5:30fetch it every day.
- 5:32Uh-huh. Look, I have a couple more
- 5:34questions about this then. Leading
- 5:35questions so people can understand it
- 5:37too. See, so you used an open API,
- 5:39meaning some source that has articles
- 5:42from the New York Times. Okay? And then
- 5:46you wrote the code and packed it all
- 5:48into the Airflow tool. Airflow is a
- 5:50thing that can just run this code by
- 5:52itself every day. We see it there at
- 5:547:00 AM, right?
- 5:55Yes. That's right.
- 5:57So, at a high level, can you describe
- 5:59the pipeline, like where you download
- 6:01from and to? Meaning,
- 6:03well, we get the data via API using
- 6:05Python. So, there’s this API, as I
- 6:08said, New York Times, and it returns
- 6:11about twenty rows of the most popular
- 6:14articles. The API requires
- 6:16authorization. And I already mentioned
- 6:17how we store the token; we can look at
- 6:19it later. By the way, if you haven't
- 6:21worked with APIs, you'll learn how to
- 6:22do it and how to use them at the
- 6:24Bootcamp. We got the data from the API.
- 6:25Let's save it. Where do we put it? Well
- 6:28, into S3. We wrap each API response
- 6:30into a parquet file and upload it to S3
- 6:32. S3 is our data lake. Yeah, let's call
- 6:35it object storage. S3 is file storage.
- 6:38Well, you can think of it as something
- 6:40like Google Drive or Dropbox. So there
- 6:41are objects, and there is a path to
- 6:42them. And if we access that path, we
- 6:44get the object. By the way, in a
- 6:46commercial context, separating the DWH,
- 6:49i.e., Greenplum, and S3 is standard
- 6:52practice. Other teams can also access
- 6:54S3, not just us. Right, well, let's say
- 6:57Airflow, we understand, it orchestrates
- 6:58, gets connections, and everything else
- 7:00. So, we accessed the API, but we can't
- 7:03load it directly, so we load it into an
- 7:05intermediate storage first. And how do
- 7:08we load it? We load parquet files.
- 7:10These are files with a columnar format.
- 7:12They are smaller in size and very fast
- 7:13due to good compression. Also, parquet
- 7:16files can store data types. Well, I
- 7:19didn't implement that here; I just
- 7:20threw the strings in as they were. We
- 7:22loaded the data into S3, and that's it;
- 7:24we store the data in raw form. So, now
- 7:27Greenplum. We transfer the data from S3
- 7:29to Greenplum. So, in S3, we have a
- 7:31bunch of unprocessed data in its raw
- 7:32form. This doesn't look at all like
- 7:34data models that you could build
- 7:36anything on. Greenplum is used here as
- 7:38a corporate data warehouse where
- 7:39cleaning, standardization, and the
- 7:41construction of those very data models
- 7:42take place. Storing such
- 7:44well-structured data in S3 simply
- 7:46doesn't make sense. Well, another team
- 7:48will need entirely different business
- 7:50logic, while in Greenplum, we will
- 7:52create business logic specifically for
- 7:53our team. First, we'll transfer the
- 7:55data to Greenplum just as it is, that
- 7:57is, we'll call this the raw data layer.
- 7:59That's it. And we'll save it there. By
- 8:01the way, how are we transferring it
- 8:03here? Well, we do it with Spark. First,
- 8:05Spark pulls the last uploaded Parquet
- 8:08file from S3, then it forms a dataframe
- 8:11and sends it to Greenplum. Let's repeat
- 8:14that one more time. So, we accessed the
- 8:16API, received the data, and uploaded it
- 8:18to S3. Now we have raw data for all
- 8:20time stored in S3. We need to process
- 8:22it somehow and build some models. For
- 8:25that, we need a DWH. We load the data
- 8:27into the DWH. Greenplum. And this is
- 8:29where those data layers begin. So, I
- 8:31already mentioned the raw layer. It's
- 8:33data as is. We take data from the raw
- 8:37layer and, using the DBT tool, build
- 8:39three more layers: stage, ODS, and in
- 8:42the raw layer, we have data that we
- 8:44simply put into the table each time.
- 8:47That is, duplicates, repeats, and
- 8:49anything else are possible there. We
- 8:51could have restarted the DAG into the
- 8:53stage layer. We only put clean data in.
- 8:56What does that mean? We take the latest
- 8:58data from the raw layer and put only
- 9:00that into the stage layer. So, in stage
- 9:02, we have it day by day, week by week,
- 9:04and everything is smooth. So, the stage
- 9:07layer serves specifically to clean this
- 9:09raw data so that the data is valid and
- 9:11there are no repeats, duplicates, etc.
- 9:14That is, we load from RAW to STG
- 9:16incrementally. By the way, we'll see
- 9:19how to do that in DBT later. In ODS, we
- 9:21then place all the data from the stage.
- 9:23Here in ODS, we cast data types,
- 9:25perhaps join them somehow, and in the
- 9:27mart layer, we store ready-made reports
- 9:29. We've formed several reports.
- 9:31Analysts and BI developers will
- 9:32interact with such reports. DBT is a
- 9:35rather interesting tool, and everything
- 9:37there is written in standard SQL. Well,
- 9:40maybe some macros are added.
- 9:41Listen, I actually have a question for
- 9:43you right here. Look, so is DBT just an
- 9:46engine, is it just code like a library,
- 9:48or is it a framework? What is it, could
- 9:52you explain? DBT is a tool for building
- 9:55data layers. So, it's a separate engine
- 9:57. We can build data layers in
- 9:59absolutely any database. Actually, dbt
- 10:02just helps us structure this, maintain
- 10:05the history of model changes, and write
- 10:07models in one place.
- 10:09So, it's basically just a wrapper for
- 10:12SQL, right? If you look at it from the
- 10:15outside, it's just SQL code wrapped in
- 10:18an additional framework like dbt.
- 10:21Yes, that's correct. We use dbt to
- 10:23implement insert logic, that is, how we
- 10:26will insert the data. So, what does a
- 10:29model actually look like? A model is
- 10:31just a standard SQL file with a query
- 10:34like "select something from" another
- 10:36model or table. Okay. In Greenplum, we
- 10:40went from a raw format, meaning
- 10:43unprocessed, to a mart layer of final
- 10:46data marts. So, our ready-made reports
- 10:48are sitting in the mart layer. Our
- 10:49analysts work in ClickHouse. And for
- 10:52those who don't know, ClickHouse is a
- 10:53database for fast analytical queries.
- 10:56It has column-oriented storage and is
- 10:58generally quite popular among analysts.
- 11:00Analysts will come here specifically to
- 11:01get their reports. And they will build
- 11:03their dashboards from here. Well,
- 11:05that's where we'll put our data from
- 11:07Greenplum. How do we move data from the
- 11:10Greenplum mart layer to ClickHouse?
- 11:12Using Spark. So, we connect Spark to
- 11:15Greenplum, take the data, create a
- 11:17dataframe, and send it to ClickHouse.
- 11:19There's an important point here, too.
- 11:21We shouldn't send the entire table. It
- 11:23would be strange to do that. We either
- 11:25have to recreate the table every time
- 11:28or, again, add data incrementally.
- 11:30We'll do it, by the way, just like we
- 11:32did from the raw layer to the DWH layer
- 11:34. That is, we'll take the data with the
- 11:37maximum date and load it. It turns out
- 11:39that we will also store only ordered
- 11:41data in ClickHouse. I actually made a
- 11:44mistake here. I was searching for the
- 11:46maximum date with Spark and then
- 11:48filtering the dataframe by it. Well,
- 11:50for those who don't know, the "max"
- 11:52function triggers a shuffle operation,
- 11:54though not a huge one. So, there are
- 11:56Spark workers, and each worker holds
- 11:58some portion of the data. We search for
- 12:01the local maximum on each worker, then
- 12:04shuffle them to find the true maximum
- 12:07among all local ones. The error isn't
- 12:10critical, but we might lose some speed.
- 12:12Okay, our final reports are stored in
- 12:16ClickHouse. Analysts use them, and they
- 12:18are happy. But as an engineer, I don't
- 12:22know how my DAG workflow is running,
- 12:24how long tasks take, or if they are
- 12:27failing. Of course, you can go into
- 12:30Airflow and check that. There’s some
- 12:33basic information on completed tasks
- 12:35there, but I decided to implement
- 12:37alerting here and put the data into
- 12:39Postgres. Actually, it’s quite simple
- 12:42to do. Airflow, as a task orchestrator,
- 12:46knows which task is running and its
- 12:48status, meaning whether it succeeded,
- 12:51failed, or retried. How much time the
- 12:54task took—all that information is
- 12:56available. So, I just took all this
- 12:58information, gathered it together, and
- 13:00sent it to a Postgres table. Actually,
- 13:03we can check the table later; it
- 13:05clearly shows what the DAG was, what
- 13:07task in the DAG it was, how long it
- 13:09took, what its status was, and when it
- 13:11was triggered. It’s very convenient.
- 13:14You can also build analytics. Well, the
- 13:16slide shows some examples of dashboards
- 13:18, yes. So, guys, does anyone have any
- 13:20questions? Unfortunately, you might
- 13:21have—
- 13:22Well, here’s a question about the
- 13:24layers, actually; let’s go ahead and
- 13:27ask it. The question is: "Could we have
- 13:31just uploaded to S3, processed it with
- 13:33Spark into a data mart, and then loaded
- 13:36it into ClickHouse, without all those
- 13:38ODS, DDS, DBT things and so on?"
- 13:42Yes, that’s actually a good question.
- 13:43With our volume of data, I think that
- 13:45would have been possible. Well, what is
- 13:47it, really? We received about 20 rows
- 13:49via API. We could have just saved them.
- 13:51Spark would have handled it very easily
- 13:53and quickly. But when the data grows,
- 13:55those layers will really be necessary.
- 13:57Without layers, your data will just
- 13:59stay in its raw form. If you perform
- 14:01any transformation, you are inevitably
- 14:03creating some kind of data layer,
- 14:04perhaps implicitly. But here, we define
- 14:06it explicitly. And we have the entire
- 14:08history of changes: how we modified
- 14:10things, what we added, what we removed.
- 14:12Maybe we added new columns, for example
- 14:14. And that is exactly what DBT helps us
- 14:16with—tracking the entire evolutionary
- 14:18history of our schema. I’ll put it
- 14:19this way. Listen, by the way, look, you
- 14:21know what question might come up? You
- 14:23talk about tracking changes. But what
- 14:25does it mean to track changes? And
- 14:26where do you track them?
- 14:27A DBT model is just a SQL file, right?
- 14:31There’s SQL code in there, but we can
- 14:33just commit it to GitHub, and we’ll
- 14:35have the full history of changes for
- 14:36that specific file. If we need to, we
- 14:38can roll back, or we can branch out if
- 14:40necessary. It’s very convenient. We
- 14:41are working exclusively with files here
- 14:43. But how are these files executed?
- 14:45Well, DBT handles that part itself. Can
- 14:47you remind me again what exactly you
- 14:48are monitoring? So, at each stage, yes,
- 14:50you can capture and see what is
- 14:52happening and record it in Postgres. Ah
- 14:54,
- 14:54yes, moreover, at each stage I can also
- 14:56exclude something from that stage. I'll
- 14:59show you later in Airflow how that's
- 15:01done. Overall, it's a fairly flexible
- 15:03system. For each task, we define how it
- 15:05will be monitored. The task has
- 15:07instructions for what to do on success
- 15:09and what to do on failure. Actually, we
- 15:12add that alerting function to both, and
- 15:15it all works. So, if our pipeline
- 15:17expands and new tasks appear, we will
- 15:19be able to easily add monitoring to
- 15:21them as well.
- 15:22For those who want code, it's time for
- 15:25some code nitpicking. Now we will take
- 15:27a look at the code itself.
- 15:28Well, let's start. So, for the API to
- 15:31S3, we have a separate task. Here it is
- 15:35. And here we specify the task ID and
- 15:40the arguments. Well, this is the URL we
- 15:43access, the folder path where we save
- 15:45it to S3, and the function. We'll get
- 15:47to that in a moment. And here, by the
- 15:49way, is that logging I mentioned. We'll
- 15:51probably look at that at the very end.
- 15:53For now, let's look at how data is
- 15:56retrieved from the API into S3. We get
- 15:59the token for the New York Times API
- 16:02and make a request. Well, I wrote a
- 16:05separate function, get_data. In it, we
- 16:07send the URL we are accessing, the
- 16:09token, and the context. And we take the
- 16:12context as the date. How do we generate
- 16:15the date? So, we send a standard GET
- 16:17request; if it's successful, we use
- 16:19pandas to do a JSON normalize. We get
- 16:22the result from the JSON first and then
- 16:25normalize it. A dataframe is formed.
- 16:28The dataframe has the following columns
- 16:30. We add the current date when Airflow
- 16:34executed or is executing this task and
- 16:37a URL hash. Each article, well,
- 16:41obviously, has some link, its own URL,
- 16:44for example, http, colon, two slashes,
- 16:47and some link to an article. But here
- 16:51we also hash this URL. Why this is
- 16:55needed, I'll explain when we talk about
- 16:58Greenplum. What is happening here? A
- 17:01regular https string turns into a
- 17:04specific number. And what's more, it
- 17:08will be the same every time if the same
- 17:10string is provided. That's all. We send
- 17:12information to the logs about how many
- 17:14rows were received. In case of failure,
- 17:16we also show what happened. We received
- 17:20the data, formed the file name,
- 17:23specified the path, the file name, and
- 17:26the date. The same one we put into the
- 17:30parquet file. If the data is not None
- 17:34and not empty, here we need to put it
- 17:38into S3. We place the finished file in
- 17:44S3 with a name, so we don't just save
- 17:48it to the machine. Well, look, we have
- 17:54Airflow already running in Docker. If
- 17:55we save the file there, well, I don't
- 17:58think it's good to just save it in the
- 18:00container. That doesn't seem reasonable
- 18:02. But using a buffer, we can put the
- 18:05file into RAM and assign a name to it.
- 18:08Actually, that's what we do. We save it
- 18:10there. And by accessing S3, this is how
- 18:14information is stored in Airflow. We
- 18:18upload our file there, specify its name
- 18:21, and specify that we're uploading from
- 18:23a buffer. We send it. If an exception
- 18:30occurs, we also log it. In any case, we
- 18:35close the buffer. It's important here
- 18:38to return the file name. Why do we
- 18:41return the file name? So that at the
- 18:45next stage, when Spark reads from S3
- 18:48and puts it into Greenplum, it knows
- 18:51the name. Where will it get it from?
- 18:55Well, Airflow has a special mechanism
- 18:57called XCom; we put our file name there
- 19:00, and the next task will take it from
- 19:03there. Okay. Let's move on. Is this
- 19:07clear? So, we accessed, formed a
- 19:10dataframe, formed a file name, uploaded
- 19:13it to S3, and returned the name of the
- 19:16file as it's stored in S3. Zhenya,
- 19:18please tell me, are there any questions
- 19:20there?
- 19:20Well, look, there's a question: can you
- 19:22just set up an API? How did you
- 19:25retrieve it? You could have just
- 19:28received the data from the response,
- 19:34and not even parse it. Just take
- 19:36exactly what you received and save it
- 19:39directly to S3, without even...
- 19:42Just the JSON,
- 19:45so that it gives you exactly that,
- 19:48without these dictionaries and such.
- 19:51Yes, of course, that would have been
- 19:52possible.
- 19:52What do you think? Might that have been
- 19:54better or worse? Because I can say that
- 19:57at the company Sravni, when we were
- 19:59there, many APIs were just like that
- 20:02too. They were slightly modified. We
- 20:05didn't always upload what we downloaded
- 20:08directly to S3. We still prepped things
- 20:11a bit, added a column here and there,
- 20:13and so on. There were some minor
- 20:15transformations after all. So I can't
- 20:17speak for the entire data engineering
- 20:20industry, saying that everyone does it
- 20:22this way or that way. Many aspects
- 20:25might just be more convenient. And I
- 20:27don't see any difference in whether it
- 20:29would be worse, or what do you think?
- 20:31Well, overall, that could have been
- 20:33done if we didn't have to create new
- 20:35features in the data. So, we actually
- 20:38needed to record the date of execution
- 20:41as well, create a new URL attribute,
- 20:43and save it exactly in that format. If
- 20:46that hadn't been required, I think, yes
- 20:49, or we could have just put all of this
- 20:52into the JSON filename. I mean, if we
- 20:55were to store the raw JSON as it
- 20:58arrives, and put the date and URL hash
- 21:00in the name, then yes, that could have
- 21:03been done. So, the data is in S3. Our
- 21:06parquet files are located in S3 at a
- 21:08specific path. Now we need to pull them
- 21:11with Spark and put them into the raw
- 21:13data layer in Greenplum. There is also
- 21:15a separate task. So, uh, we get the
- 21:18Spark session. Well, I will actually
- 21:21demonstrate this. Uh, we just needed to
- 21:23use it several times, so it was decided
- 21:25to write it out as a separate function.
- 21:27So, Spark is launched with certain
- 21:30parameters. Uh, well, actually, to be
- 21:32honest, we will need all of this here.
- 21:34I see MinIO here, that's our S3. I see
- 21:37Hive, we need that too. Postgres, yes,
- 21:40most likely everything is needed here.
- 21:44So, we got the Spark session, got Spark
- 21:46, and we get the filename. As I said
- 21:49before, we previously stored the
- 21:51filename in XCom, which is the path in
- 21:53S3. And now we pull it from XCom. Here,
- 21:58the task ID is the one that was
- 22:01supposed to store that path there. Now
- 22:05our Spark path is stored here. So, get
- 22:09Spark DF. Also a separate function. We
- 22:12pass in the Spark session and the
- 22:14filename. Let's see what it does. Aha,
- 22:19it just goes to S3, accesses the path,
- 22:23reads Parquet, creates a dataframe, and
- 22:26returns it. In case of failure, it
- 22:29displays it in the console. Okay, so
- 22:32that's fine. The Spark dataframe is
- 22:34here. Now we need to put it into
- 22:37Greenplum. There is a separate function
- 22:41where we pass the Greenplum table name
- 22:44and our Spark dataframe. Let's see how
- 22:47it does that. So, we get a connection
- 22:50to Greenplum. We create a...GP URL,
- 22:56meaning we access Greenplum via JDBC
- 23:01and write the data there. Actually,
- 23:05data write. Data is the...Spark
- 23:09dataframe that we want to write. Write
- 23:12method. We access it via the path we
- 23:16formed. We write the data. And which
- 23:19table do we specify? Well, the login
- 23:23and password are for Greenplum, which
- 23:25we fetched right here. And we inform
- 23:28that yes, everything is loaded. Well,
- 23:30if something went wrong at this stage,
- 23:32we will get an error. That's it. So, we
- 23:34did it, something didn't work out, and
- 23:36again, we display that. Well, I don't
- 23:39know what could possibly go wrong here,
- 23:40so I’m just outputting the exception
- 23:42as it is to the console, and that’s
- 23:43it. Well, in the end, naturally, we
- 23:46stop the Spark session that we obtained
- 23:49here. I’d prefer for each separate
- 23:51task to start its own session and shut
- 23:53it down. So, let me summarize. We took
- 23:57the file path from XCom, created a
- 24:00dataframe, and sent it to Greenplum. So
- 24:04, dbt is triggered in Airflow using
- 24:07this operator, the dbt run operator.
- 24:11But the way it’s set up in our
- 24:14bootcamp is that we place our dbt
- 24:17models in specific paths, and then when
- 24:20we run a GitHub Action deploy, they are
- 24:24moved to other folders. It all starts
- 24:27with the models. Each dbt run executes
- 24:31a certain model: stage, ODS, mart. Yes,
- 24:36we actually have several models in the
- 24:38mart, but I only upload one to
- 24:40ClickHouse. I don't need all of them. I
- 24:43think we should take a look at the
- 24:45models now. So, this is our stage model
- 24:48. Let me remind you, we take from RAW,
- 24:53from the raw table, and load data
- 24:56incrementally into the stage. How is
- 24:59this done? Well, I select the fields
- 25:01that I need. I specify the source.
- 25:05I’ve moved it into a separate file.
- 25:07You can take a look here. The schema
- 25:10and the table. These are all the data
- 25:13we have in Greenplum. I'm referencing
- 25:17it. And if it's incremental, which it
- 25:21is, I check if the date—I select the
- 25:25rows where the date is greater than the
- 25:29maximum date. What does this mean? This
- 25:36is our current layer; so in the stage,
- 25:40I look at the raw table for rows where
- 25:45the date is greater than the max date
- 25:49in the stage table. I hope that's clear
- 25:55. I don't need all the data from RAW. I
- 25:59only take the data that hasn't been
- 26:01recorded yet. But I have specified—or
- 26:05rather, described—this model here,
- 26:07giving it a description and a config.
- 26:10Which one? Materializing the table as
- 26:13incremental. That is, in fact, why this
- 26:15flag worked. If we hadn't specified
- 26:18this, it wouldn't have worked.
- 26:21Materialization. And the schema, by the
- 26:25way, I didn't tell you. Let me remind
- 26:30you that when we were fetching data via
- 26:33the API and putting it into S3, between
- 26:36those two processes, I created two new
- 26:39features: date and URL hash. Why do we
- 26:43need it? Why do we need this number? We
- 26:47are loading the data into Greenplum.
- 26:50Greenplum can be thought of as a
- 26:52multi-shard Postgres, meaning several
- 26:55Postgres instances combined into an MPP
- 26:58architecture. So there is some, well,
- 27:01master coordinator, and it tells them
- 27:03what to do. Our data is stored on
- 27:07different shards, on different Postgres
- 27:11instances. How do we distribute this
- 27:13data between these Postgres instances?
- 27:16The "distributed by" clause answers
- 27:18exactly this question. We distribute
- 27:21our data by key, and depending on what
- 27:24the URL hash is, it gets sent to a
- 27:27specific Postgres instance. Actually,
- 27:31distribution, I think, is mandatory in
- 27:34Greenplum. Zhenya, can you remind me?
- 27:37Well, distribution needs to be spread
- 27:39out, yes, absolutely. The only question
- 27:41is, by which field?
- 27:42Yes, exactly, yes. It's precisely for
- 27:46this field that our hash is quite
- 27:48high-cardinality, meaning it has many
- 27:50unique values, which is why we spread
- 27:53it out based on it. Okay. That's how
- 27:58the data ended up in our staging layer.
- 28:00Meaning we got it incrementally from
- 28:03raw and stored it as is. Moving on. ODS
- 28:06. Things are much simpler here. We take
- 28:09everything entirely from the staging
- 28:11layer. We take the table. We access the
- 28:14columns, but we convert the data types
- 28:17into a normal format. So here we format
- 28:20the dates. We don't need to convert the
- 28:25URL hash, but it was originally stored
- 28:27in Parquet as a number, so nothing
- 28:29needs to be done here. Of course, ODS,
- 28:33in theory, should be richer. But, I'll
- 28:37emphasize again, our data volume is
- 28:39very small. There is nothing, so to
- 28:43speak, to transform here. Okay. And now
- 28:46we go to the marts. In marts, we made
- 28:50one model by simply taking everything
- 28:53from ODS. We made another model as a
- 28:56report, meaning we grouped by update
- 28:59date and section. The section in the
- 29:02API is simply what kind of article it
- 29:05is, what it belongs to. So we grouped
- 29:08them and looked at how many articles we
- 29:10have for each such section. Here we
- 29:13also grouped and looked by article
- 29:16author, how many articles belong to one
- 29:18author on that same day, the download
- 29:21date, so to speak. And the same thing
- 29:24with sections. We won't take these
- 29:26reports to ClickHouse; they aren't
- 29:27needed there. But we will take this one
- 29:30.
- 29:31I can anticipate a question: why choose
- 29:33"append" as the incremental strategy?
- 29:36Because there are more incremental
- 29:38strategies in dbt. Well, the question
- 29:42is, why that one? What do we need to do
- 29:46? Just put the new data at the end, or,
- 29:49you know, just add new data. We aren't
- 29:53going to overwrite the past, so,
- 29:54effectively, "append." Well, here each
- 29:59task runs our models and we execute
- 30:05them sequentially, first staging, then
- 30:08marts. We combine them into a group.
- 30:11The group is needed solely for a nice
- 30:15visual. I think I showed the graph. It
- 30:18makes it clear that this is a dbt group
- 30:21, and the picture becomes more
- 30:23intuitive. Here it is. So, the DBT
- 30:27group moved data from the RA table to
- 30:30stage, from stage to DS, and from UDS
- 30:32to the mart. The ready reports are
- 30:35stored in the mart. Now we need to take
- 30:38those ready reports and put them into
- 30:40ClickHouse. There is a separate task
- 30:42for this. Actually, here it is. Again,
- 30:49we take a Spark session. What do we
- 30:52need? We need to take it from Greenplum
- 30:54and put it into ClickHouse. Naturally,
- 30:55we need connections to these two
- 30:57databases. We create them. Okay. We
- 31:03read from the mart layer. So, there is
- 31:07this schema in Greenplum. This is the
- 31:10table stored there. We read it, and it
- 31:13becomes our Spark dataframe. So, by the
- 31:18way, I mentioned this error in the
- 31:20presentation. Look at what I’m doing.
- 31:24I’ve retrieved the entire volume of
- 31:28data, find the maximum, then filter by
- 31:31that maximum, keeping only the most
- 31:35recent data in this dataframe. In my
- 31:39opinion, doing it this way wasn't the
- 31:40best approach. But it works, and the
- 31:43error isn't that critical. I think if
- 31:46there’s more data in the mart layer,
- 31:48then yes, slowdowns might actually
- 31:50occur. So, I’ve got the fresh data
- 31:54and am writing it to ClickHouse. Here
- 31:57is the distributed table. There is
- 31:59another complexity here. Our bootcamp
- 32:03ClickHouse is multi-shard. Actually, in
- 32:09the scenario, in the very first group
- 32:12of Airflow tasks, you can see how local
- 32:16tables are created on each ClickHouse
- 32:19shard. And how a distributed table is
- 32:23then created, which links these local
- 32:25ones; it acts as a reference to them, a
- 32:28unified interface, you could say. We
- 32:31write the data specifically into this
- 32:33distributed table, and then it gets
- 32:35distributed across those ClickHouse
- 32:37shards. Overall, I think if it were
- 32:39single-shard, it would be simpler, but
- 32:41this is how we have it. We wrote these
- 32:44fresh data points to the ClickHouse
- 32:47table, informed the user, and finalized
- 32:50by stopping the Spark session. I think
- 32:53that was the last task. Ah, yes, in the
- 32:56DAG, the data has flowed into
- 32:59ClickHouse.
- 33:00And there’s this question; I didn’t
- 33:02quite understand what it was about.
- 33:04What we don't carry over to ClickHouse,
- 33:06why is it there? Where is it used?
- 33:08Ah, yes, by the way, that’s a good
- 33:09question. I see why we leave the entire
- 33:13mart table in Greenplum if not
- 33:15everything goes to ClickHouse, only the
- 33:18latest. Yes, that is indeed how it
- 33:21happens. In Greenplum, in our mart
- 33:24layer, we have almost the same report
- 33:26as in ClickHouse, we just trim what’s
- 33:28needed and deliver to ClickHouse via a
- 33:30simple insert. Well, let me answer it
- 33:35this way: maybe our mart report in
- 33:37Greenplum is needed by someone else,
- 33:39another team, and not just analysts,
- 33:41and they would like to see this final
- 33:43report in Greenplum and use it, maybe
- 33:46process it in their own way. So, our
- 33:49report might not just be for ClickHouse
- 33:51analysts; our report might be needed by
- 33:53someone in Greenplum as well. Well, let
- 33:55it stay there.
- 33:57Listen, here's a question: did you know
- 34:00all of this? For example, DBT, Airflow,
- 34:03Spark, ClickHouse, Greenplum, did you
- 34:05know these, or was it introductory
- 34:09knowledge?
- 34:10Yes. What backend, well, you know, what
- 34:12was your background? Right. Well, I
- 34:15knew Postgres, I knew a little
- 34:18ClickHouse, I knew a bit of Spark and
- 34:22used Airflow, but it was, I think, I
- 34:25didn't write as well in Airflow back
- 34:28then as I do now, and I also didn't
- 34:32know ClickHouse at all, especially if
- 34:35it's multi-sharded. So, I heard about
- 34:39local and distributed tables for the
- 34:41first time here. Well, Spark, I don't
- 34:43think so, my level definitely went up.
- 34:46My Spark tasks at work were much
- 34:48simpler. Of course, there are some
- 34:50tasks here at the Bootcamp that are not
- 34:52similar, but at work, everything was
- 34:54simplified. What did I learn here?
- 34:55Greenplum. I never touched it before.
- 34:58DBT. I'll put it this way, data models
- 35:01existed. Without them, I think it's
- 35:02impossible to build anything. Data
- 35:04models were there, but without DBT.
- 35:06Someone would write something like: "
- 35:08Why do we need DBT if we can just use
- 35:09SQL?""Yes, indeed, that's how it was
- 35:11for us. But I liked DBT more. It has
- 35:15more capabilities than just storing
- 35:18separate SQL files. At work, I had at
- 35:21most two tasks, maybe three, but
- 35:23nothing this large. I had never written
- 35:25anything like this.
- 35:26Listen, a question for the future. What
- 35:28do you think, your DAG is quite large,
- 35:30isn't it? How many lines? 500 lines,
- 35:32yeah, maybe more. Can your DBT models
- 35:35be built, let's say, remotely? Yes, you
- 35:39don't have any DBT model code written
- 35:42directly here, but your Spark API,
- 35:45especially Spark tasks, and possibly
- 35:47other tasks, are written right in your
- 35:50DAG. How do you feel about moving all
- 35:54of that into separate scripts so that
- 35:56your DAG serves only as a description
- 35:59of the tasks and jobs?
- 36:01Yes, of course. So, do the same as with
- 36:04dbt, yeah, write them into separate
- 36:06Python files and just call them. Yes, I
- 36:09think we can. In general, I think
- 36:10that's how people do it. Nobody is
- 36:12going to write 500 lines of code.
- 36:15Okay, let me try to ask the question.
- 36:18Oh, and you try to answer. Go ahead,
- 36:20watch. Hello everyone. I know this is a
- 36:23pet project. What in theory could be
- 36:25unnecessary here? Which technologies
- 36:27are needed here, and which are not?
- 36:29What can be improved, or conversely,
- 36:30removed? I'm just starting to learn
- 36:32everything myself. Well, actually, I
- 36:34wanted to ask this question too. What
- 36:36do you think? It's clear that it's a
- 36:37pet project, and you could probably do
- 36:40everything here with just one tool,
- 36:41maybe Python, and set it all up on cron
- 36:43, right? But what do you think? Maybe
- 36:48it's good for the future, like for
- 36:50bigger data, or what could be removed,
- 36:53what could we do without?
- 36:56Okay, well, let's see. If we are saving
- 36:58this volume of data, say an API returns
- 37:00a little bit every day, how would I do
- 37:02it? So there is an API, it returns
- 37:06something, and we want to see it in
- 37:07ClickHouse, right? Well, what is the
- 37:10task? Get it from the API into
- 37:13ClickHouse so analysts can use it. What
- 37:16would I do? I would also put the data
- 37:21into S3, just so it's there, have a
- 37:23history of that data, and be able to
- 37:25re-access it. And in ClickHouse, I
- 37:28would create an integration table with
- 37:31an integration engine for S3. And I
- 37:34would just access all this data through
- 37:37the integration table first, then load
- 37:39it into a regular ClickHouse table with
- 37:42a standard engine, like MergeTree.
- 37:46Actually, since we didn't really change
- 37:48the data here, and there were no
- 37:51special transformations, we would just
- 37:53take it as is from S3 and put it into
- 37:55ClickHouse using the S3 integration
- 37:58engine. By the way, I think we have a
- 38:00project like that. Yes.
- 38:02Well, yeah, we do.
- 38:03And so, Greenplum and dbt definitely
- 38:05wouldn't be needed here.
- 38:07But tell me, let's say this is a fairly
- 38:09complex pet project—well, not complex
- 38:11, but for a beginner, in my opinion, it
- 38:13would be difficult. And regarding
- 38:15priorities, what would you rank first?
- 38:17What you definitely need to know, what
- 38:20you definitely need to figure out, and
- 38:23what should be left for last in your
- 38:25studies? Well, if a person wants to
- 38:27build a pet project like this.
- 38:28Let's start from the very beginning.
- 38:31First, I think it's the databases, you
- 38:34know, knowing Postgres or ClickHouse.
- 38:38That was my starting point, those were
- 38:40my initial inputs. I knew Postgres, I
- 38:42knew a bit of Click, so it made sense
- 38:44to me. Why is it important to know?
- 38:46Well, you have to store data somehow.
- 38:48How are we going to store and keep
- 38:50anything if we don't know SQL queries
- 38:52or how different databases work? Moving
- 38:55data from one database to another
- 38:57requires some kind of ETL tool. Spark.
- 39:02I think I might even put Spark in first
- 39:04place, but where would I connect and
- 39:06where would I save if I don't know
- 39:08databases? So, databases come first,
- 39:11then Spark. If you know these, I think
- 39:15a good half of all data engineering
- 39:18will be covered.
- 39:20Uh-huh. Then Airflow, probably, too.
- 39:22Well, you need to write it somewhere so
- 39:24it runs on a schedule.
- 39:25You could use Airflow, or you could use
- 39:27Cron, as we mentioned. At work, I've
- 39:30used Cron before.
- 39:32So, DBT comes at the very end, is that
- 39:34right? In terms of priority?
- 39:36Yes? Well, you can live without it. And
- 39:38, by the way, it wouldn't be best
- 39:41practice, but you could build models
- 39:45with Spark instead of DBT. You could
- 39:48create stage or ODS schemas and use
- 39:52Spark to move one table to another,
- 39:55transforming them along the way. You
- 39:58could do that, yes; it's slow,
- 40:01yes, it's moving data over the network,
- 40:05which isn't ideal. Usually, what DBT
- 40:08does under the hood is an" insert into
- 40:12"table," select all from "another table
- 40:15, which is very fast. But Spark
- 40:17involves moving data over the network.
- 40:19Well, you can do it with Spark too,
- 40:21sure. If there isn't much data, I think
- 40:24it would be a great solution. Basically
- 40:27, the architectural pipeline is exactly
- 40:31the same as when Volodya and I worked
- 40:34at Sravni. It's not something invented
- 40:37from scratch, reinventing the wheel.
- 40:41All these technologies and combinations
- 40:43, as I said at the beginning, have been
- 40:45around for a long time. The main thing
- 40:48is to keep them in mind, not get
- 40:50confused, and set them up reasonably
- 40:52well. By the way, I really want to add
- 40:55these combinations and examples to the
- 40:57roadmap. We have a lecture on company
- 40:58architecture in our bootcamp, something
- 41:00like that. Also, let's add the bundles,
- 41:02because that's really cool. I mean, you
- 41:04can clearly see how things can be done,
- 41:06how they are done, and how everyone is
- 41:08used to them, which is fine. Basically,
- 41:10the best practices, so we don't have to
- 41:12reinvent the wheel. Oh, Artyom, a
- 41:13question: how long did it take you to
- 41:15build this whole thing? How much time
- 41:16did you spend on it? Well, let's say,
- 41:19solid 3 to 4 days.
- 41:22How often did you use ChatGPT, for what
- 41:24, and where did it help you? How did it
- 41:26speed things up exactly? Or maybe you
- 41:28read the documentation or used a book?
- 41:30Yes, by the way, that's a great
- 41:31question. I quite often had problems
- 41:34with connectors, like how to get the
- 41:37connection itself, how to parse it, you
- 41:40know, get the host, port, all that
- 41:42stuff. I definitely asked about that.
- 41:46By the way, I didn't know how to create
- 41:47a TaskGroup. I had seen it, I had seen
- 41:51how it worked, but specifically how to
- 41:53build one. Well, it turned out to be
- 41:55quite simple. I mean, let's say ChatGPT
- 41:57mostly provided me with examples.
- 42:00Tell me, did you use ChatGPT as a
- 42:02separate chat or within your code
- 42:04editor as an agent?
- 42:06No, separate, I just asked questions in
- 42:08the chat.
- 42:08Alright, I suggest we wrap things up
- 42:10then. Jen, thanks everyone for
- 42:11attending our stream. We'll try to do
- 42:14this for every cohort, maybe, well,
- 42:16depending on the interest, of course.
- 42:19Come to our bootcamp. This absolutely
- 42:22had to be mentioned here. And in
- 42:24general, there will be a recording on
- 42:26YouTube. Go there, like, and subscribe
- 42:28to the channel to help us grow.
- 42:30Yes, thank you very much, Artyom, for
- 42:32the project. It's honestly an
- 42:33incredibly huge amount of work done. At
- 42:35the same time, what we shared, I think,
- 42:36was clear to everyone. Guys, as for the
- 42:39bootcamp, the cohort starts on July 1st
- 42:41. As they say, we are always in the
- 42:43channel, we always talk about it.
- 42:45There's a bootcamp bot there, go in,
- 42:46type" I want to join the bootcamp "if
- 42:48you're really interested. And regarding
- 42:50the code and the presentation, we'll
- 42:52prepare and package both. Possibly in a
- 42:55separate repository or in a post, we'll
- 42:58see what's more convenient. It will
- 43:01also be available along with the
- 43:02recording. So, you'll just have to wait
- 43:04a little while, and it will be
- 43:05available completely for free. And it's
- 43:08great, like it's usually, like a
- 43:09roadmap. Artyom, say a few words if you
- 43:11have anything to say, something to
- 43:12share, maybe, if you want to. Really
- 43:15anything, yeah. Guys, thank you very
- 43:17much for inviting me, thank you very
- 43:18much for letting me show all of this.
- 43:20And a big thank you to the chat as well
- 43:22. I see a lot of solid advice here. I
- 43:25will definitely take everything into
- 43:26account, learn everything, and put it
- 43:28all into practice, guys. Thank you.
- 43:31Alright, that's it, thanks to everyone
- 43:32who was here. Guys, I'm ending the
- 43:34stream, I've stopped it. That's all,
- 43:35bye everyone. y
About this transcript
This page contains the full transcript of Пет Проект для Data Engineer | Инженер Данных by Евгений Виндюков, generated from the public captions YouTube serves with the video. The transcript has 7,091 words across 1,038 segments, with the original timestamps preserved so you can click any line to jump to that moment in the embedded player.
What you can do with it
Use the transcript to take notes, quote the speaker, build a study guide, generate a summary with ChatGPT or Claude via the YouTube Summary tool, or export it as a timed subtitle file with YouTube to SRT. You can also re-open it in the transcriber to translate the transcript into 100+ languages.
Free YouTube transcript tool
YouTube2Text is a free YouTube transcript generator — no signup, no daily limit. Paste any YouTube link and get the full transcript instantly, with timestamps, click-to-jump, translation to 100+ languages, AI prompts for ChatGPT, Claude, and Gemini, and exports to TXT, SRT, VTT, or Markdown.