Repository files navigation

Stream Processing Project

This is a personal project that I started with the goal to study and use some of the most used stream-processing technologies.

I am not interested in testing the scalability, robustness or any infrastructure-related theme on any of these tools. Also, I am not instered in advanced features of each technology. The purpose here is to understand how would you code a simple stream processing pipeline (something not too far from a Hello World for streaming pipelines).

Technologies Used

Architecture

Architecture

Using the Python library faker, we create fake transaction data to simulate a generic e-commerce. Our main goal is to aggregate those transactions in 3-second windows to know our income in real-time (as we shall see below).

The data is produced by the fake_data.py script. One can notice that we hard-coded a list of 10 user_ids. This was done to help debugging and ensure correctness of our stream pipelines, since each user_id start with a number from 0 to 9 and the each user pays a fixed amount equal to ten times its leading number (with 0 being interpreted as 10). As an example, the user with leading number 5 will always pay the amount of 50.

We then use both Apache Flink, kSQL and Apache Spark to process the data. Our goal is to group all the paid transactions into 3 second windows and calculate:

  • the sum of what was paid;
  • the number of transactions;
  • the average of transactions.

Instructions

In order to run the project, just run the command make project. This will start all the docker containers and generate (almost) all the necessary resources in each container.

Since Flink doesn't have a REST API to programatically define the stream transformations, one must run

make flink-sql

to start a container with the Flink SQL Client.

Flink SQL Client

Once in there, just copy & paste all the SQL instructions located in the flink folder (in the same order as the files are numbered). By the end of it, the CLI should report that a Job has been submited to Flink:

Flink Job Created

You can go to localhost:8081 and you shall see the Job Running:

Flink Job Running

With all that done, run make data N={{ number }} to send {{ number }} of messages to Kafka. The script will send the {{ number }} messages using 10 parallel processes.

Sending Data using N=10

You can visualize the messages sent using the Kafka UI available at localhost:8080. The messages are sent to the transactions topic. The aggregations are stored at the transactions_aggregate_{{ technology }} topic.

Messages in Kafka

Then, you can see the resulting tables for results coming from kSQL in pinot using the pinot-controller interface at localhost:9000.

Data in Pinot

If needed, more commands are available in the Makefile.

Notes

  • In case it is needed, the ksql CLI is available by running make ksql-cli.

Learnings

  • We need to send a key to Kafka topic in order to use table in kSQL. See Topic #6 in https://www.confluent.io/blog/troubleshooting-ksql-part-1/.
  • The ROWTIME pseudo-column available in kSQL: https://docs.ksqldb.io/en/latest/developer-guide/ksqldb-reference/create-stream/#rowtime
  • Apache Pinot does a incredible job of compacting the Kafka records that are "wrong" (e.g. without the complete calculations from the window, example below) coming from kSQL (and it is fair to imagine that it would do the same job for the data coming from Spark). One could imagine what are the performance gains (or losses) if we set Kafkas log.cleanup.policy to compact instead of the default delete. Obviously this approach would need defining a key for each message, but a natural candidate for this would be the window_start (or a hash of it).

Ongoing calculations for some window

Design Choices

We delibery chose to not use default Docker images of said technologies whenever possible. Since each one can be downloaded and run locally, we chose to do (almost) the same thing using Docker containers. This allowed us to know more about the configurations available in each one (as one can see in the config) folder.

References

About

Repository for a personal project that uses the most common stream-processing technologies

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

Repository files navigation

Stream Processing Project

This is a personal project that I started with the goal to study and use some of the most used stream-processing technologies.

I am not interested in testing the scalability, robustness or any infrastructure-related theme on any of these tools. Also, I am not instered in advanced features of each technology. The purpose here is to understand how would you code a simple stream processing pipeline (something not too far from a Hello World for streaming pipelines).

Technologies Used

Architecture

Architecture

Using the Python library faker, we create fake transaction data to simulate a generic e-commerce. Our main goal is to aggregate those transactions in 3-second windows to know our income in real-time (as we shall see below).

The data is produced by the fake_data.py script. One can notice that we hard-coded a list of 10 user_ids. This was done to help debugging and ensure correctness of our stream pipelines, since each user_id start with a number from 0 to 9 and the each user pays a fixed amount equal to ten times its leading number (with 0 being interpreted as 10). As an example, the user with leading number 5 will always pay the amount of 50.

We then use both Apache Flink, kSQL and Apache Spark to process the data. Our goal is to group all the paid transactions into 3 second windows and calculate:

  • the sum of what was paid;
  • the number of transactions;
  • the average of transactions.

Instructions

In order to run the project, just run the command make project. This will start all the docker containers and generate (almost) all the necessary resources in each container.

Since Flink doesn't have a REST API to programatically define the stream transformations, one must run

make flink-sql

to start a container with the Flink SQL Client.

Flink SQL Client

Once in there, just copy & paste all the SQL instructions located in the flink folder (in the same order as the files are numbered). By the end of it, the CLI should report that a Job has been submited to Flink:

Flink Job Created

You can go to localhost:8081 and you shall see the Job Running:

Flink Job Running

With all that done, run make data N={{ number }} to send {{ number }} of messages to Kafka. The script will send the {{ number }} messages using 10 parallel processes.

Sending Data using N=10

You can visualize the messages sent using the Kafka UI available at localhost:8080. The messages are sent to the transactions topic. The aggregations are stored at the transactions_aggregate_{{ technology }} topic.

Messages in Kafka

Then, you can see the resulting tables for results coming from kSQL in pinot using the pinot-controller interface at localhost:9000.

Data in Pinot

If needed, more commands are available in the Makefile.

Notes

  • In case it is needed, the ksql CLI is available by running make ksql-cli.

Learnings

  • We need to send a key to Kafka topic in order to use table in kSQL. See Topic #6 in https://www.confluent.io/blog/troubleshooting-ksql-part-1/.
  • The ROWTIME pseudo-column available in kSQL: https://docs.ksqldb.io/en/latest/developer-guide/ksqldb-reference/create-stream/#rowtime
  • Apache Pinot does a incredible job of compacting the Kafka records that are "wrong" (e.g. without the complete calculations from the window, example below) coming from kSQL (and it is fair to imagine that it would do the same job for the data coming from Spark). One could imagine what are the performance gains (or losses) if we set Kafkas log.cleanup.policy to compact instead of the default delete. Obviously this approach would need defining a key for each message, but a natural candidate for this would be the window_start (or a hash of it).

Ongoing calculations for some window

Design Choices

We delibery chose to not use default Docker images of said technologies whenever possible. Since each one can be downloaded and run locally, we chose to do (almost) the same thing using Docker containers. This allowed us to know more about the configurations available in each one (as one can see in the config) folder.

References

About

Repository for a personal project that uses the most common stream-processing technologies

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Repository files navigation

Stream Processing Project

This is a personal project that I started with the goal to study and use some of the most used stream-processing technologies.

I am not interested in testing the scalability, robustness or any infrastructure-related theme on any of these tools. Also, I am not instered in advanced features of each technology. The purpose here is to understand how would you code a simple stream processing pipeline (something not too far from a Hello World for streaming pipelines).

Technologies Used

Architecture

Architecture

Using the Python library faker, we create fake transaction data to simulate a generic e-commerce. Our main goal is to aggregate those transactions in 3-second windows to know our income in real-time (as we shall see below).

The data is produced by the fake_data.py script. One can notice that we hard-coded a list of 10 user_ids. This was done to help debugging and ensure correctness of our stream pipelines, since each user_id start with a number from 0 to 9 and the each user pays a fixed amount equal to ten times its leading number (with 0 being interpreted as 10). As an example, the user with leading number 5 will always pay the amount of 50.

We then use both Apache Flink, kSQL and Apache Spark to process the data. Our goal is to group all the paid transactions into 3 second windows and calculate:

  • the sum of what was paid;
  • the number of transactions;
  • the average of transactions.

Instructions

In order to run the project, just run the command make project. This will start all the docker containers and generate (almost) all the necessary resources in each container.

Since Flink doesn't have a REST API to programatically define the stream transformations, one must run

make flink-sql

to start a container with the Flink SQL Client.

Flink SQL Client

Once in there, just copy & paste all the SQL instructions located in the flink folder (in the same order as the files are numbered). By the end of it, the CLI should report that a Job has been submited to Flink:

Flink Job Created

You can go to localhost:8081 and you shall see the Job Running:

Flink Job Running

With all that done, run make data N={{ number }} to send {{ number }} of messages to Kafka. The script will send the {{ number }} messages using 10 parallel processes.

Sending Data using N=10

You can visualize the messages sent using the Kafka UI available at localhost:8080. The messages are sent to the transactions topic. The aggregations are stored at the transactions_aggregate_{{ technology }} topic.

Messages in Kafka

Then, you can see the resulting tables for results coming from kSQL in pinot using the pinot-controller interface at localhost:9000.

Data in Pinot

If needed, more commands are available in the Makefile.

Notes

  • In case it is needed, the ksql CLI is available by running make ksql-cli.

Learnings

  • We need to send a key to Kafka topic in order to use table in kSQL. See Topic #6 in https://www.confluent.io/blog/troubleshooting-ksql-part-1/.
  • The ROWTIME pseudo-column available in kSQL: https://docs.ksqldb.io/en/latest/developer-guide/ksqldb-reference/create-stream/#rowtime
  • Apache Pinot does a incredible job of compacting the Kafka records that are "wrong" (e.g. without the complete calculations from the window, example below) coming from kSQL (and it is fair to imagine that it would do the same job for the data coming from Spark). One could imagine what are the performance gains (or losses) if we set Kafkas log.cleanup.policy to compact instead of the default delete. Obviously this approach would need defining a key for each message, but a natural candidate for this would be the window_start (or a hash of it).

Ongoing calculations for some window

Design Choices

We delibery chose to not use default Docker images of said technologies whenever possible. Since each one can be downloaded and run locally, we chose to do (almost) the same thing using Docker containers. This allowed us to know more about the configurations available in each one (as one can see in the config) folder.

References

About

Repository for a personal project that uses the most common stream-processing technologies

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Repository files navigation

Stream Processing Project

This is a personal project that I started with the goal to study and use some of the most used stream-processing technologies.

I am not interested in testing the scalability, robustness or any infrastructure-related theme on any of these tools. Also, I am not instered in advanced features of each technology. The purpose here is to understand how would you code a simple stream processing pipeline (something not too far from a Hello World for streaming pipelines).

Technologies Used

Architecture

Architecture

Using the Python library faker, we create fake transaction data to simulate a generic e-commerce. Our main goal is to aggregate those transactions in 3-second windows to know our income in real-time (as we shall see below).

The data is produced by the fake_data.py script. One can notice that we hard-coded a list of 10 user_ids. This was done to help debugging and ensure correctness of our stream pipelines, since each user_id start with a number from 0 to 9 and the each user pays a fixed amount equal to ten times its leading number (with 0 being interpreted as 10). As an example, the user with leading number 5 will always pay the amount of 50.

We then use both Apache Flink, kSQL and Apache Spark to process the data. Our goal is to group all the paid transactions into 3 second windows and calculate:

  • the sum of what was paid;
  • the number of transactions;
  • the average of transactions.

Instructions

In order to run the project, just run the command make project. This will start all the docker containers and generate (almost) all the necessary resources in each container.

Since Flink doesn't have a REST API to programatically define the stream transformations, one must run

make flink-sql

to start a container with the Flink SQL Client.

Flink SQL Client

Once in there, just copy & paste all the SQL instructions located in the flink folder (in the same order as the files are numbered). By the end of it, the CLI should report that a Job has been submited to Flink:

Flink Job Created

You can go to localhost:8081 and you shall see the Job Running:

Flink Job Running

With all that done, run make data N={{ number }} to send {{ number }} of messages to Kafka. The script will send the {{ number }} messages using 10 parallel processes.

Sending Data using N=10

You can visualize the messages sent using the Kafka UI available at localhost:8080. The messages are sent to the transactions topic. The aggregations are stored at the transactions_aggregate_{{ technology }} topic.

Messages in Kafka

Then, you can see the resulting tables for results coming from kSQL in pinot using the pinot-controller interface at localhost:9000.

Data in Pinot

If needed, more commands are available in the Makefile.

Notes

  • In case it is needed, the ksql CLI is available by running make ksql-cli.

Learnings

  • We need to send a key to Kafka topic in order to use table in kSQL. See Topic #6 in https://www.confluent.io/blog/troubleshooting-ksql-part-1/.
  • The ROWTIME pseudo-column available in kSQL: https://docs.ksqldb.io/en/latest/developer-guide/ksqldb-reference/create-stream/#rowtime
  • Apache Pinot does a incredible job of compacting the Kafka records that are "wrong" (e.g. without the complete calculations from the window, example below) coming from kSQL (and it is fair to imagine that it would do the same job for the data coming from Spark). One could imagine what are the performance gains (or losses) if we set Kafkas log.cleanup.policy to compact instead of the default delete. Obviously this approach would need defining a key for each message, but a natural candidate for this would be the window_start (or a hash of it).

Ongoing calculations for some window

Design Choices

We delibery chose to not use default Docker images of said technologies whenever possible. Since each one can be downloaded and run locally, we chose to do (almost) the same thing using Docker containers. This allowed us to know more about the configurations available in each one (as one can see in the config) folder.

References

About

Repository for a personal project that uses the most common stream-processing technologies

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

Repository files navigation

Stream Processing Project

This is a personal project that I started with the goal to study and use some of the most used stream-processing technologies.

I am not interested in testing the scalability, robustness or any infrastructure-related theme on any of these tools. Also, I am not instered in advanced features of each technology. The purpose here is to understand how would you code a simple stream processing pipeline (something not too far from a Hello World for streaming pipelines).

Technologies Used

Architecture

Architecture

Using the Python library faker, we create fake transaction data to simulate a generic e-commerce. Our main goal is to aggregate those transactions in 3-second windows to know our income in real-time (as we shall see below).

The data is produced by the fake_data.py script. One can notice that we hard-coded a list of 10 user_ids. This was done to help debugging and ensure correctness of our stream pipelines, since each user_id start with a number from 0 to 9 and the each user pays a fixed amount equal to ten times its leading number (with 0 being interpreted as 10). As an example, the user with leading number 5 will always pay the amount of 50.

We then use both Apache Flink, kSQL and Apache Spark to process the data. Our goal is to group all the paid transactions into 3 second windows and calculate:

  • the sum of what was paid;
  • the number of transactions;
  • the average of transactions.

Instructions

In order to run the project, just run the command make project. This will start all the docker containers and generate (almost) all the necessary resources in each container.

Since Flink doesn't have a REST API to programatically define the stream transformations, one must run

make flink-sql

to start a container with the Flink SQL Client.

Flink SQL Client

Once in there, just copy & paste all the SQL instructions located in the flink folder (in the same order as the files are numbered). By the end of it, the CLI should report that a Job has been submited to Flink:

Flink Job Created

You can go to localhost:8081 and you shall see the Job Running:

Flink Job Running

With all that done, run make data N={{ number }} to send {{ number }} of messages to Kafka. The script will send the {{ number }} messages using 10 parallel processes.

Sending Data using N=10

You can visualize the messages sent using the Kafka UI available at localhost:8080. The messages are sent to the transactions topic. The aggregations are stored at the transactions_aggregate_{{ technology }} topic.

Messages in Kafka

Then, you can see the resulting tables for results coming from kSQL in pinot using the pinot-controller interface at localhost:9000.

Data in Pinot

If needed, more commands are available in the Makefile.

Notes

  • In case it is needed, the ksql CLI is available by running make ksql-cli.

Learnings

  • We need to send a key to Kafka topic in order to use table in kSQL. See Topic #6 in https://www.confluent.io/blog/troubleshooting-ksql-part-1/.
  • The ROWTIME pseudo-column available in kSQL: https://docs.ksqldb.io/en/latest/developer-guide/ksqldb-reference/create-stream/#rowtime
  • Apache Pinot does a incredible job of compacting the Kafka records that are "wrong" (e.g. without the complete calculations from the window, example below) coming from kSQL (and it is fair to imagine that it would do the same job for the data coming from Spark). One could imagine what are the performance gains (or losses) if we set Kafkas log.cleanup.policy to compact instead of the default delete. Obviously this approach would need defining a key for each message, but a natural candidate for this would be the window_start (or a hash of it).

Ongoing calculations for some window

Design Choices

We delibery chose to not use default Docker images of said technologies whenever possible. Since each one can be downloaded and run locally, we chose to do (almost) the same thing using Docker containers. This allowed us to know more about the configurations available in each one (as one can see in the config) folder.

References

About

Repository for a personal project that uses the most common stream-processing technologies

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Repository files navigation

Stream Processing Project

This is a personal project that I started with the goal to study and use some of the most used stream-processing technologies.

I am not interested in testing the scalability, robustness or any infrastructure-related theme on any of these tools. Also, I am not instered in advanced features of each technology. The purpose here is to understand how would you code a simple stream processing pipeline (something not too far from a Hello World for streaming pipelines).

Technologies Used

Architecture

Architecture

Using the Python library faker, we create fake transaction data to simulate a generic e-commerce. Our main goal is to aggregate those transactions in 3-second windows to know our income in real-time (as we shall see below).

The data is produced by the fake_data.py script. One can notice that we hard-coded a list of 10 user_ids. This was done to help debugging and ensure correctness of our stream pipelines, since each user_id start with a number from 0 to 9 and the each user pays a fixed amount equal to ten times its leading number (with 0 being interpreted as 10). As an example, the user with leading number 5 will always pay the amount of 50.

We then use both Apache Flink, kSQL and Apache Spark to process the data. Our goal is to group all the paid transactions into 3 second windows and calculate:

  • the sum of what was paid;
  • the number of transactions;
  • the average of transactions.

Instructions

In order to run the project, just run the command make project. This will start all the docker containers and generate (almost) all the necessary resources in each container.

Since Flink doesn't have a REST API to programatically define the stream transformations, one must run

make flink-sql

to start a container with the Flink SQL Client.

Flink SQL Client

Once in there, just copy & paste all the SQL instructions located in the flink folder (in the same order as the files are numbered). By the end of it, the CLI should report that a Job has been submited to Flink:

Flink Job Created

You can go to localhost:8081 and you shall see the Job Running:

Flink Job Running

With all that done, run make data N={{ number }} to send {{ number }} of messages to Kafka. The script will send the {{ number }} messages using 10 parallel processes.

Sending Data using N=10

You can visualize the messages sent using the Kafka UI available at localhost:8080. The messages are sent to the transactions topic. The aggregations are stored at the transactions_aggregate_{{ technology }} topic.

Messages in Kafka

Then, you can see the resulting tables for results coming from kSQL in pinot using the pinot-controller interface at localhost:9000.

Data in Pinot

If needed, more commands are available in the Makefile.

Notes

  • In case it is needed, the ksql CLI is available by running make ksql-cli.

Learnings

  • We need to send a key to Kafka topic in order to use table in kSQL. See Topic #6 in https://www.confluent.io/blog/troubleshooting-ksql-part-1/.
  • The ROWTIME pseudo-column available in kSQL: https://docs.ksqldb.io/en/latest/developer-guide/ksqldb-reference/create-stream/#rowtime
  • Apache Pinot does a incredible job of compacting the Kafka records that are "wrong" (e.g. without the complete calculations from the window, example below) coming from kSQL (and it is fair to imagine that it would do the same job for the data coming from Spark). One could imagine what are the performance gains (or losses) if we set Kafkas log.cleanup.policy to compact instead of the default delete. Obviously this approach would need defining a key for each message, but a natural candidate for this would be the window_start (or a hash of it).

Ongoing calculations for some window

Design Choices

We delibery chose to not use default Docker images of said technologies whenever possible. Since each one can be downloaded and run locally, we chose to do (almost) the same thing using Docker containers. This allowed us to know more about the configurations available in each one (as one can see in the config) folder.

References

About

Repository for a personal project that uses the most common stream-processing technologies

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Repository files navigation

Stream Processing Project

This is a personal project that I started with the goal to study and use some of the most used stream-processing technologies.

I am not interested in testing the scalability, robustness or any infrastructure-related theme on any of these tools. Also, I am not instered in advanced features of each technology. The purpose here is to understand how would you code a simple stream processing pipeline (something not too far from a Hello World for streaming pipelines).

Technologies Used

Architecture

Architecture

Using the Python library faker, we create fake transaction data to simulate a generic e-commerce. Our main goal is to aggregate those transactions in 3-second windows to know our income in real-time (as we shall see below).

The data is produced by the fake_data.py script. One can notice that we hard-coded a list of 10 user_ids. This was done to help debugging and ensure correctness of our stream pipelines, since each user_id start with a number from 0 to 9 and the each user pays a fixed amount equal to ten times its leading number (with 0 being interpreted as 10). As an example, the user with leading number 5 will always pay the amount of 50.

We then use both Apache Flink, kSQL and Apache Spark to process the data. Our goal is to group all the paid transactions into 3 second windows and calculate:

  • the sum of what was paid;
  • the number of transactions;
  • the average of transactions.

Instructions

In order to run the project, just run the command make project. This will start all the docker containers and generate (almost) all the necessary resources in each container.

Since Flink doesn't have a REST API to programatically define the stream transformations, one must run

make flink-sql

to start a container with the Flink SQL Client.

Flink SQL Client

Once in there, just copy & paste all the SQL instructions located in the flink folder (in the same order as the files are numbered). By the end of it, the CLI should report that a Job has been submited to Flink:

Flink Job Created

You can go to localhost:8081 and you shall see the Job Running:

Flink Job Running

With all that done, run make data N={{ number }} to send {{ number }} of messages to Kafka. The script will send the {{ number }} messages using 10 parallel processes.

Sending Data using N=10

You can visualize the messages sent using the Kafka UI available at localhost:8080. The messages are sent to the transactions topic. The aggregations are stored at the transactions_aggregate_{{ technology }} topic.

Messages in Kafka

Then, you can see the resulting tables for results coming from kSQL in pinot using the pinot-controller interface at localhost:9000.

Data in Pinot

If needed, more commands are available in the Makefile.

Notes

  • In case it is needed, the ksql CLI is available by running make ksql-cli.

Learnings

  • We need to send a key to Kafka topic in order to use table in kSQL. See Topic #6 in https://www.confluent.io/blog/troubleshooting-ksql-part-1/.
  • The ROWTIME pseudo-column available in kSQL: https://docs.ksqldb.io/en/latest/developer-guide/ksqldb-reference/create-stream/#rowtime
  • Apache Pinot does a incredible job of compacting the Kafka records that are "wrong" (e.g. without the complete calculations from the window, example below) coming from kSQL (and it is fair to imagine that it would do the same job for the data coming from Spark). One could imagine what are the performance gains (or losses) if we set Kafkas log.cleanup.policy to compact instead of the default delete. Obviously this approach would need defining a key for each message, but a natural candidate for this would be the window_start (or a hash of it).

Ongoing calculations for some window

Design Choices

We delibery chose to not use default Docker images of said technologies whenever possible. Since each one can be downloaded and run locally, we chose to do (almost) the same thing using Docker containers. This allowed us to know more about the configurations available in each one (as one can see in the config) folder.

References

About

Repository for a personal project that uses the most common stream-processing technologies

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

Repository files navigation

Stream Processing Project

This is a personal project that I started with the goal to study and use some of the most used stream-processing technologies.

I am not interested in testing the scalability, robustness or any infrastructure-related theme on any of these tools. Also, I am not instered in advanced features of each technology. The purpose here is to understand how would you code a simple stream processing pipeline (something not too far from a Hello World for streaming pipelines).

Technologies Used

Architecture

Architecture

Using the Python library faker, we create fake transaction data to simulate a generic e-commerce. Our main goal is to aggregate those transactions in 3-second windows to know our income in real-time (as we shall see below).

The data is produced by the fake_data.py script. One can notice that we hard-coded a list of 10 user_ids. This was done to help debugging and ensure correctness of our stream pipelines, since each user_id start with a number from 0 to 9 and the each user pays a fixed amount equal to ten times its leading number (with 0 being interpreted as 10). As an example, the user with leading number 5 will always pay the amount of 50.

We then use both Apache Flink, kSQL and Apache Spark to process the data. Our goal is to group all the paid transactions into 3 second windows and calculate:

  • the sum of what was paid;
  • the number of transactions;
  • the average of transactions.

Instructions

In order to run the project, just run the command make project. This will start all the docker containers and generate (almost) all the necessary resources in each container.

Since Flink doesn't have a REST API to programatically define the stream transformations, one must run

make flink-sql

to start a container with the Flink SQL Client.

Flink SQL Client

Once in there, just copy & paste all the SQL instructions located in the flink folder (in the same order as the files are numbered). By the end of it, the CLI should report that a Job has been submited to Flink:

Flink Job Created

You can go to localhost:8081 and you shall see the Job Running:

Flink Job Running

With all that done, run make data N={{ number }} to send {{ number }} of messages to Kafka. The script will send the {{ number }} messages using 10 parallel processes.

Sending Data using N=10

You can visualize the messages sent using the Kafka UI available at localhost:8080. The messages are sent to the transactions topic. The aggregations are stored at the transactions_aggregate_{{ technology }} topic.

Messages in Kafka

Then, you can see the resulting tables for results coming from kSQL in pinot using the pinot-controller interface at localhost:9000.

Data in Pinot

If needed, more commands are available in the Makefile.

Notes

  • In case it is needed, the ksql CLI is available by running make ksql-cli.

Learnings

  • We need to send a key to Kafka topic in order to use table in kSQL. See Topic #6 in https://www.confluent.io/blog/troubleshooting-ksql-part-1/.
  • The ROWTIME pseudo-column available in kSQL: https://docs.ksqldb.io/en/latest/developer-guide/ksqldb-reference/create-stream/#rowtime
  • Apache Pinot does a incredible job of compacting the Kafka records that are "wrong" (e.g. without the complete calculations from the window, example below) coming from kSQL (and it is fair to imagine that it would do the same job for the data coming from Spark). One could imagine what are the performance gains (or losses) if we set Kafkas log.cleanup.policy to compact instead of the default delete. Obviously this approach would need defining a key for each message, but a natural candidate for this would be the window_start (or a hash of it).

Ongoing calculations for some window

Design Choices

We delibery chose to not use default Docker images of said technologies whenever possible. Since each one can be downloaded and run locally, we chose to do (almost) the same thing using Docker containers. This allowed us to know more about the configurations available in each one (as one can see in the config) folder.

References

About

Repository for a personal project that uses the most common stream-processing technologies

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

Languages