Repository files navigation

Spark Internals

Spark Version: 1.0.2
Doc Version: 1.0.2.0

Authors

Weibo/Twitter IDNameContributions
@JerryLeadLijie XuAuthor of the original Chinese version, and English version update
@juhanlolHan JUEnglish version and update (Chapter 0, 1, 3, 4, and 7)
@invkrhHao RenEnglish version and update (Chapter 2, 5, and 6)

Introduction

This series discuss the design and implementation of Apache Spark, with focuses on its design principles, execution mechanisms, system architecture and performance optimization. In addition, there's some comparisons with Hadoop MapReduce in terms of design and implementation. I'm reluctant to call this document a "code walkthrough", because the goal is not to analyze each piece of code in the project, but to understand the whole system in a systematic way (through analyzing the execution procedure of a Spark job, from its creation to completion).

There're many ways to discuss a computer system. Here, We've chosen a problem-driven approach. Firstly one concret problem is introduced, then it gets analyzed step by step. We'll start from a typical Spark example job and then discuss all the related important system modules. I believe that this approach is better than diving into each module right from the beginning.

The target audiences of this series are geeks who want to have a deeper understanding of Apache Spark as well as other distributed computing frameworks.

I'll try my best to keep this documentation up to date with Spark since it's a fast evolving project with an active community. The documentation's main version is in sync with Spark's version. The additional number at the end represents the documentation's update version.

For more academic oriented discussion, please check out Matei's PHD thesis and other related papers. You can also have a look at my blog (in Chinese) blog.

I haven't been writing such complete documentation for a while. Last time it was about three years ago when I'm studying Andrew Ng's ML course. I was really motivated at that time! This time I've spent 20+ days on this document, from the summer break till now (August 2014). Most of the time is spent on debugging, drawing diagrams and thinking how to put my ideas in the right way. I hope you find this series helpful.

Contents

We start from the creation of a Spark job, and then discuss its execution. Finally, we dive into some related system modules and features.

  1. Overview Overview of Apache Spark
  2. Job logical plan Logical plan of a job (data dependency graph)
  3. Job physical plan Physical plan
  4. Shuffle details Shuffle process
  5. Architecture Coordination of system modules in job execution
  6. Cache and Checkpoint Cache and Checkpoint
  7. Broadcast Broadcast feature
  8. Job Scheduling TODO
  9. Fault-tolerance TODO

Chinses Verison is at markdown/.

The documentation is written in markdown. The pdf version is also available here.

If you're under Max OS X, I recommand MacDown with a github theme for reading.

Gitbook (Chinese version)

Thanks @Yourtion for creating the gitbook version.

Online reading http://spark-internals.books.yourtion.com/

Downloads

Examples

I've created some examples to debug the system during the writing, they are avaible under SparkLearning/src/internals.

Acknowledgement

I appreciate the help from the following in providing solutions and ideas for some detailed issues:

  • @Andrew-Xia Participated in the discussion of BlockManager's implemetation's impact on broadcast(rdd).

  • @CrazyJVM Participated in the discussion of BlockManager's implementation.

  • @王联辉 Participated in the discussion of BlockManager's implementation.

Thanks to the following for complementing the document:

Weibo IDChapterContentRevision status
@OopsOutOfMemoryOverviewRelation between workers and executors and Summary on Spark Executor Driver's Resouce Management (in Chinese)There's not yet a conclusion on this subject since its implementation is still changing, a link to the blog is added

Thanks to the following for finding errors:

Weibo IdChapterError/IssueRevision status
@JoshuawangzjOverviewWhen multiple applications are running, multiple Backend process will be createdCorrected, but need to be confirmed. No idea on how to control the number of Backend processes
@_cs_cmOverviewLatest groupByKey() has removed the mapValues() operation, there's no MapValuesRDD generatedFixed groupByKey() related diagrams and text
@染染生起JobLogicalPlanN:N relation in FullDepedency N:N is a NarrowDependencyModified the description of NarrowDependency into 3 different cases with detaild explaination, clearer than the 2 cases explaination before
@zzl0Fisrt four chaptersLots of typos,such as "groupByKey has generated the 3 following RDDs",should be 2. Check pull requestAll fixed
@左手牵右手TELCache and Broadcast chapterLots of typosAll fixed
@cloud-fanJobLogicalPlanSome arrows in the Cogroup() diagram should be colored redAll fixed
@CrazyJvmShuffle detailsStarting from Spark 1.1, the default value for spark.shuffle.file.buffer.kb is 32k, not 100kAll fixed

Special thanks to @明风Andy for his great support.

Special thanks to the rockers (including researchers, developers and users) who participate in the design, implementation and discussion of big data systems.

About

Notes talking about the design and implementation of Apache Spark

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

, '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

Spark Internals

Spark Version: 1.0.2
Doc Version: 1.0.2.0

Authors

Weibo/Twitter IDNameContributions
@JerryLeadLijie XuAuthor of the original Chinese version, and English version update
@juhanlolHan JUEnglish version and update (Chapter 0, 1, 3, 4, and 7)
@invkrhHao RenEnglish version and update (Chapter 2, 5, and 6)

Introduction

This series discuss the design and implementation of Apache Spark, with focuses on its design principles, execution mechanisms, system architecture and performance optimization. In addition, there's some comparisons with Hadoop MapReduce in terms of design and implementation. I'm reluctant to call this document a "code walkthrough", because the goal is not to analyze each piece of code in the project, but to understand the whole system in a systematic way (through analyzing the execution procedure of a Spark job, from its creation to completion).

There're many ways to discuss a computer system. Here, We've chosen a problem-driven approach. Firstly one concret problem is introduced, then it gets analyzed step by step. We'll start from a typical Spark example job and then discuss all the related important system modules. I believe that this approach is better than diving into each module right from the beginning.

The target audiences of this series are geeks who want to have a deeper understanding of Apache Spark as well as other distributed computing frameworks.

I'll try my best to keep this documentation up to date with Spark since it's a fast evolving project with an active community. The documentation's main version is in sync with Spark's version. The additional number at the end represents the documentation's update version.

For more academic oriented discussion, please check out Matei's PHD thesis and other related papers. You can also have a look at my blog (in Chinese) blog.

I haven't been writing such complete documentation for a while. Last time it was about three years ago when I'm studying Andrew Ng's ML course. I was really motivated at that time! This time I've spent 20+ days on this document, from the summer break till now (August 2014). Most of the time is spent on debugging, drawing diagrams and thinking how to put my ideas in the right way. I hope you find this series helpful.

Contents

We start from the creation of a Spark job, and then discuss its execution. Finally, we dive into some related system modules and features.

  1. Overview Overview of Apache Spark
  2. Job logical plan Logical plan of a job (data dependency graph)
  3. Job physical plan Physical plan
  4. Shuffle details Shuffle process
  5. Architecture Coordination of system modules in job execution
  6. Cache and Checkpoint Cache and Checkpoint
  7. Broadcast Broadcast feature
  8. Job Scheduling TODO
  9. Fault-tolerance TODO

Chinses Verison is at markdown/.

The documentation is written in markdown. The pdf version is also available here.

If you're under Max OS X, I recommand MacDown with a github theme for reading.

Gitbook (Chinese version)

Thanks @Yourtion for creating the gitbook version.

Online reading http://spark-internals.books.yourtion.com/

Downloads

Examples

I've created some examples to debug the system during the writing, they are avaible under SparkLearning/src/internals.

Acknowledgement

I appreciate the help from the following in providing solutions and ideas for some detailed issues:

  • @Andrew-Xia Participated in the discussion of BlockManager's implemetation's impact on broadcast(rdd).

  • @CrazyJVM Participated in the discussion of BlockManager's implementation.

  • @王联辉 Participated in the discussion of BlockManager's implementation.

Thanks to the following for complementing the document:

Weibo IDChapterContentRevision status
@OopsOutOfMemoryOverviewRelation between workers and executors and Summary on Spark Executor Driver's Resouce Management (in Chinese)There's not yet a conclusion on this subject since its implementation is still changing, a link to the blog is added

Thanks to the following for finding errors:

Weibo IdChapterError/IssueRevision status
@JoshuawangzjOverviewWhen multiple applications are running, multiple Backend process will be createdCorrected, but need to be confirmed. No idea on how to control the number of Backend processes
@_cs_cmOverviewLatest groupByKey() has removed the mapValues() operation, there's no MapValuesRDD generatedFixed groupByKey() related diagrams and text
@染染生起JobLogicalPlanN:N relation in FullDepedency N:N is a NarrowDependencyModified the description of NarrowDependency into 3 different cases with detaild explaination, clearer than the 2 cases explaination before
@zzl0Fisrt four chaptersLots of typos,such as "groupByKey has generated the 3 following RDDs",should be 2. Check pull requestAll fixed
@左手牵右手TELCache and Broadcast chapterLots of typosAll fixed
@cloud-fanJobLogicalPlanSome arrows in the Cogroup() diagram should be colored redAll fixed
@CrazyJvmShuffle detailsStarting from Spark 1.1, the default value for spark.shuffle.file.buffer.kb is 32k, not 100kAll fixed

Special thanks to @明风Andy for his great support.

Special thanks to the rockers (including researchers, developers and users) who participate in the design, implementation and discussion of big data systems.

About

Notes talking about the design and implementation of Apache Spark

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

, '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

Spark Internals

Spark Version: 1.0.2
Doc Version: 1.0.2.0

Authors

Weibo/Twitter IDNameContributions
@JerryLeadLijie XuAuthor of the original Chinese version, and English version update
@juhanlolHan JUEnglish version and update (Chapter 0, 1, 3, 4, and 7)
@invkrhHao RenEnglish version and update (Chapter 2, 5, and 6)

Introduction

This series discuss the design and implementation of Apache Spark, with focuses on its design principles, execution mechanisms, system architecture and performance optimization. In addition, there's some comparisons with Hadoop MapReduce in terms of design and implementation. I'm reluctant to call this document a "code walkthrough", because the goal is not to analyze each piece of code in the project, but to understand the whole system in a systematic way (through analyzing the execution procedure of a Spark job, from its creation to completion).

There're many ways to discuss a computer system. Here, We've chosen a problem-driven approach. Firstly one concret problem is introduced, then it gets analyzed step by step. We'll start from a typical Spark example job and then discuss all the related important system modules. I believe that this approach is better than diving into each module right from the beginning.

The target audiences of this series are geeks who want to have a deeper understanding of Apache Spark as well as other distributed computing frameworks.

I'll try my best to keep this documentation up to date with Spark since it's a fast evolving project with an active community. The documentation's main version is in sync with Spark's version. The additional number at the end represents the documentation's update version.

For more academic oriented discussion, please check out Matei's PHD thesis and other related papers. You can also have a look at my blog (in Chinese) blog.

I haven't been writing such complete documentation for a while. Last time it was about three years ago when I'm studying Andrew Ng's ML course. I was really motivated at that time! This time I've spent 20+ days on this document, from the summer break till now (August 2014). Most of the time is spent on debugging, drawing diagrams and thinking how to put my ideas in the right way. I hope you find this series helpful.

Contents

We start from the creation of a Spark job, and then discuss its execution. Finally, we dive into some related system modules and features.

  1. Overview Overview of Apache Spark
  2. Job logical plan Logical plan of a job (data dependency graph)
  3. Job physical plan Physical plan
  4. Shuffle details Shuffle process
  5. Architecture Coordination of system modules in job execution
  6. Cache and Checkpoint Cache and Checkpoint
  7. Broadcast Broadcast feature
  8. Job Scheduling TODO
  9. Fault-tolerance TODO

Chinses Verison is at markdown/.

The documentation is written in markdown. The pdf version is also available here.

If you're under Max OS X, I recommand MacDown with a github theme for reading.

Gitbook (Chinese version)

Thanks @Yourtion for creating the gitbook version.

Online reading http://spark-internals.books.yourtion.com/

Downloads

Examples

I've created some examples to debug the system during the writing, they are avaible under SparkLearning/src/internals.

Acknowledgement

I appreciate the help from the following in providing solutions and ideas for some detailed issues:

  • @Andrew-Xia Participated in the discussion of BlockManager's implemetation's impact on broadcast(rdd).

  • @CrazyJVM Participated in the discussion of BlockManager's implementation.

  • @王联辉 Participated in the discussion of BlockManager's implementation.

Thanks to the following for complementing the document:

Weibo IDChapterContentRevision status
@OopsOutOfMemoryOverviewRelation between workers and executors and Summary on Spark Executor Driver's Resouce Management (in Chinese)There's not yet a conclusion on this subject since its implementation is still changing, a link to the blog is added

Thanks to the following for finding errors:

Weibo IdChapterError/IssueRevision status
@JoshuawangzjOverviewWhen multiple applications are running, multiple Backend process will be createdCorrected, but need to be confirmed. No idea on how to control the number of Backend processes
@_cs_cmOverviewLatest groupByKey() has removed the mapValues() operation, there's no MapValuesRDD generatedFixed groupByKey() related diagrams and text
@染染生起JobLogicalPlanN:N relation in FullDepedency N:N is a NarrowDependencyModified the description of NarrowDependency into 3 different cases with detaild explaination, clearer than the 2 cases explaination before
@zzl0Fisrt four chaptersLots of typos,such as "groupByKey has generated the 3 following RDDs",should be 2. Check pull requestAll fixed
@左手牵右手TELCache and Broadcast chapterLots of typosAll fixed
@cloud-fanJobLogicalPlanSome arrows in the Cogroup() diagram should be colored redAll fixed
@CrazyJvmShuffle detailsStarting from Spark 1.1, the default value for spark.shuffle.file.buffer.kb is 32k, not 100kAll fixed

Special thanks to @明风Andy for his great support.

Special thanks to the rockers (including researchers, developers and users) who participate in the design, implementation and discussion of big data systems.

About

Notes talking about the design and implementation of Apache Spark

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

, '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

Spark Internals

Spark Version: 1.0.2
Doc Version: 1.0.2.0

Authors

Weibo/Twitter IDNameContributions
@JerryLeadLijie XuAuthor of the original Chinese version, and English version update
@juhanlolHan JUEnglish version and update (Chapter 0, 1, 3, 4, and 7)
@invkrhHao RenEnglish version and update (Chapter 2, 5, and 6)

Introduction

This series discuss the design and implementation of Apache Spark, with focuses on its design principles, execution mechanisms, system architecture and performance optimization. In addition, there's some comparisons with Hadoop MapReduce in terms of design and implementation. I'm reluctant to call this document a "code walkthrough", because the goal is not to analyze each piece of code in the project, but to understand the whole system in a systematic way (through analyzing the execution procedure of a Spark job, from its creation to completion).

There're many ways to discuss a computer system. Here, We've chosen a problem-driven approach. Firstly one concret problem is introduced, then it gets analyzed step by step. We'll start from a typical Spark example job and then discuss all the related important system modules. I believe that this approach is better than diving into each module right from the beginning.

The target audiences of this series are geeks who want to have a deeper understanding of Apache Spark as well as other distributed computing frameworks.

I'll try my best to keep this documentation up to date with Spark since it's a fast evolving project with an active community. The documentation's main version is in sync with Spark's version. The additional number at the end represents the documentation's update version.

For more academic oriented discussion, please check out Matei's PHD thesis and other related papers. You can also have a look at my blog (in Chinese) blog.

I haven't been writing such complete documentation for a while. Last time it was about three years ago when I'm studying Andrew Ng's ML course. I was really motivated at that time! This time I've spent 20+ days on this document, from the summer break till now (August 2014). Most of the time is spent on debugging, drawing diagrams and thinking how to put my ideas in the right way. I hope you find this series helpful.

Contents

We start from the creation of a Spark job, and then discuss its execution. Finally, we dive into some related system modules and features.

  1. Overview Overview of Apache Spark
  2. Job logical plan Logical plan of a job (data dependency graph)
  3. Job physical plan Physical plan
  4. Shuffle details Shuffle process
  5. Architecture Coordination of system modules in job execution
  6. Cache and Checkpoint Cache and Checkpoint
  7. Broadcast Broadcast feature
  8. Job Scheduling TODO
  9. Fault-tolerance TODO

Chinses Verison is at markdown/.

The documentation is written in markdown. The pdf version is also available here.

If you're under Max OS X, I recommand MacDown with a github theme for reading.

Gitbook (Chinese version)

Thanks @Yourtion for creating the gitbook version.

Online reading http://spark-internals.books.yourtion.com/

Downloads

Examples

I've created some examples to debug the system during the writing, they are avaible under SparkLearning/src/internals.

Acknowledgement

I appreciate the help from the following in providing solutions and ideas for some detailed issues:

  • @Andrew-Xia Participated in the discussion of BlockManager's implemetation's impact on broadcast(rdd).

  • @CrazyJVM Participated in the discussion of BlockManager's implementation.

  • @王联辉 Participated in the discussion of BlockManager's implementation.

Thanks to the following for complementing the document:

Weibo IDChapterContentRevision status
@OopsOutOfMemoryOverviewRelation between workers and executors and Summary on Spark Executor Driver's Resouce Management (in Chinese)There's not yet a conclusion on this subject since its implementation is still changing, a link to the blog is added

Thanks to the following for finding errors:

Weibo IdChapterError/IssueRevision status
@JoshuawangzjOverviewWhen multiple applications are running, multiple Backend process will be createdCorrected, but need to be confirmed. No idea on how to control the number of Backend processes
@_cs_cmOverviewLatest groupByKey() has removed the mapValues() operation, there's no MapValuesRDD generatedFixed groupByKey() related diagrams and text
@染染生起JobLogicalPlanN:N relation in FullDepedency N:N is a NarrowDependencyModified the description of NarrowDependency into 3 different cases with detaild explaination, clearer than the 2 cases explaination before
@zzl0Fisrt four chaptersLots of typos,such as "groupByKey has generated the 3 following RDDs",should be 2. Check pull requestAll fixed
@左手牵右手TELCache and Broadcast chapterLots of typosAll fixed
@cloud-fanJobLogicalPlanSome arrows in the Cogroup() diagram should be colored redAll fixed
@CrazyJvmShuffle detailsStarting from Spark 1.1, the default value for spark.shuffle.file.buffer.kb is 32k, not 100kAll fixed

Special thanks to @明风Andy for his great support.

Special thanks to the rockers (including researchers, developers and users) who participate in the design, implementation and discussion of big data systems.

About

Notes talking about the design and implementation of Apache Spark

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

, '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

Spark Internals

Spark Version: 1.0.2
Doc Version: 1.0.2.0

Authors

Weibo/Twitter IDNameContributions
@JerryLeadLijie XuAuthor of the original Chinese version, and English version update
@juhanlolHan JUEnglish version and update (Chapter 0, 1, 3, 4, and 7)
@invkrhHao RenEnglish version and update (Chapter 2, 5, and 6)

Introduction

This series discuss the design and implementation of Apache Spark, with focuses on its design principles, execution mechanisms, system architecture and performance optimization. In addition, there's some comparisons with Hadoop MapReduce in terms of design and implementation. I'm reluctant to call this document a "code walkthrough", because the goal is not to analyze each piece of code in the project, but to understand the whole system in a systematic way (through analyzing the execution procedure of a Spark job, from its creation to completion).

There're many ways to discuss a computer system. Here, We've chosen a problem-driven approach. Firstly one concret problem is introduced, then it gets analyzed step by step. We'll start from a typical Spark example job and then discuss all the related important system modules. I believe that this approach is better than diving into each module right from the beginning.

The target audiences of this series are geeks who want to have a deeper understanding of Apache Spark as well as other distributed computing frameworks.

I'll try my best to keep this documentation up to date with Spark since it's a fast evolving project with an active community. The documentation's main version is in sync with Spark's version. The additional number at the end represents the documentation's update version.

For more academic oriented discussion, please check out Matei's PHD thesis and other related papers. You can also have a look at my blog (in Chinese) blog.

I haven't been writing such complete documentation for a while. Last time it was about three years ago when I'm studying Andrew Ng's ML course. I was really motivated at that time! This time I've spent 20+ days on this document, from the summer break till now (August 2014). Most of the time is spent on debugging, drawing diagrams and thinking how to put my ideas in the right way. I hope you find this series helpful.

Contents

We start from the creation of a Spark job, and then discuss its execution. Finally, we dive into some related system modules and features.

  1. Overview Overview of Apache Spark
  2. Job logical plan Logical plan of a job (data dependency graph)
  3. Job physical plan Physical plan
  4. Shuffle details Shuffle process
  5. Architecture Coordination of system modules in job execution
  6. Cache and Checkpoint Cache and Checkpoint
  7. Broadcast Broadcast feature
  8. Job Scheduling TODO
  9. Fault-tolerance TODO

Chinses Verison is at markdown/.

The documentation is written in markdown. The pdf version is also available here.

If you're under Max OS X, I recommand MacDown with a github theme for reading.

Gitbook (Chinese version)

Thanks @Yourtion for creating the gitbook version.

Online reading http://spark-internals.books.yourtion.com/

Downloads

Examples

I've created some examples to debug the system during the writing, they are avaible under SparkLearning/src/internals.

Acknowledgement

I appreciate the help from the following in providing solutions and ideas for some detailed issues:

  • @Andrew-Xia Participated in the discussion of BlockManager's implemetation's impact on broadcast(rdd).

  • @CrazyJVM Participated in the discussion of BlockManager's implementation.

  • @王联辉 Participated in the discussion of BlockManager's implementation.

Thanks to the following for complementing the document:

Weibo IDChapterContentRevision status
@OopsOutOfMemoryOverviewRelation between workers and executors and Summary on Spark Executor Driver's Resouce Management (in Chinese)There's not yet a conclusion on this subject since its implementation is still changing, a link to the blog is added

Thanks to the following for finding errors:

Weibo IdChapterError/IssueRevision status
@JoshuawangzjOverviewWhen multiple applications are running, multiple Backend process will be createdCorrected, but need to be confirmed. No idea on how to control the number of Backend processes
@_cs_cmOverviewLatest groupByKey() has removed the mapValues() operation, there's no MapValuesRDD generatedFixed groupByKey() related diagrams and text
@染染生起JobLogicalPlanN:N relation in FullDepedency N:N is a NarrowDependencyModified the description of NarrowDependency into 3 different cases with detaild explaination, clearer than the 2 cases explaination before
@zzl0Fisrt four chaptersLots of typos,such as "groupByKey has generated the 3 following RDDs",should be 2. Check pull requestAll fixed
@左手牵右手TELCache and Broadcast chapterLots of typosAll fixed
@cloud-fanJobLogicalPlanSome arrows in the Cogroup() diagram should be colored redAll fixed
@CrazyJvmShuffle detailsStarting from Spark 1.1, the default value for spark.shuffle.file.buffer.kb is 32k, not 100kAll fixed

Special thanks to @明风Andy for his great support.

Special thanks to the rockers (including researchers, developers and users) who participate in the design, implementation and discussion of big data systems.

About

Notes talking about the design and implementation of Apache Spark

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

, '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

Spark Internals

Spark Version: 1.0.2
Doc Version: 1.0.2.0

Authors

Weibo/Twitter IDNameContributions
@JerryLeadLijie XuAuthor of the original Chinese version, and English version update
@juhanlolHan JUEnglish version and update (Chapter 0, 1, 3, 4, and 7)
@invkrhHao RenEnglish version and update (Chapter 2, 5, and 6)

Introduction

This series discuss the design and implementation of Apache Spark, with focuses on its design principles, execution mechanisms, system architecture and performance optimization. In addition, there's some comparisons with Hadoop MapReduce in terms of design and implementation. I'm reluctant to call this document a "code walkthrough", because the goal is not to analyze each piece of code in the project, but to understand the whole system in a systematic way (through analyzing the execution procedure of a Spark job, from its creation to completion).

There're many ways to discuss a computer system. Here, We've chosen a problem-driven approach. Firstly one concret problem is introduced, then it gets analyzed step by step. We'll start from a typical Spark example job and then discuss all the related important system modules. I believe that this approach is better than diving into each module right from the beginning.

The target audiences of this series are geeks who want to have a deeper understanding of Apache Spark as well as other distributed computing frameworks.

I'll try my best to keep this documentation up to date with Spark since it's a fast evolving project with an active community. The documentation's main version is in sync with Spark's version. The additional number at the end represents the documentation's update version.

For more academic oriented discussion, please check out Matei's PHD thesis and other related papers. You can also have a look at my blog (in Chinese) blog.

I haven't been writing such complete documentation for a while. Last time it was about three years ago when I'm studying Andrew Ng's ML course. I was really motivated at that time! This time I've spent 20+ days on this document, from the summer break till now (August 2014). Most of the time is spent on debugging, drawing diagrams and thinking how to put my ideas in the right way. I hope you find this series helpful.

Contents

We start from the creation of a Spark job, and then discuss its execution. Finally, we dive into some related system modules and features.

  1. Overview Overview of Apache Spark
  2. Job logical plan Logical plan of a job (data dependency graph)
  3. Job physical plan Physical plan
  4. Shuffle details Shuffle process
  5. Architecture Coordination of system modules in job execution
  6. Cache and Checkpoint Cache and Checkpoint
  7. Broadcast Broadcast feature
  8. Job Scheduling TODO
  9. Fault-tolerance TODO

Chinses Verison is at markdown/.

The documentation is written in markdown. The pdf version is also available here.

If you're under Max OS X, I recommand MacDown with a github theme for reading.

Gitbook (Chinese version)

Thanks @Yourtion for creating the gitbook version.

Online reading http://spark-internals.books.yourtion.com/

Downloads

Examples

I've created some examples to debug the system during the writing, they are avaible under SparkLearning/src/internals.

Acknowledgement

I appreciate the help from the following in providing solutions and ideas for some detailed issues:

  • @Andrew-Xia Participated in the discussion of BlockManager's implemetation's impact on broadcast(rdd).

  • @CrazyJVM Participated in the discussion of BlockManager's implementation.

  • @王联辉 Participated in the discussion of BlockManager's implementation.

Thanks to the following for complementing the document:

Weibo IDChapterContentRevision status
@OopsOutOfMemoryOverviewRelation between workers and executors and Summary on Spark Executor Driver's Resouce Management (in Chinese)There's not yet a conclusion on this subject since its implementation is still changing, a link to the blog is added

Thanks to the following for finding errors:

Weibo IdChapterError/IssueRevision status
@JoshuawangzjOverviewWhen multiple applications are running, multiple Backend process will be createdCorrected, but need to be confirmed. No idea on how to control the number of Backend processes
@_cs_cmOverviewLatest groupByKey() has removed the mapValues() operation, there's no MapValuesRDD generatedFixed groupByKey() related diagrams and text
@染染生起JobLogicalPlanN:N relation in FullDepedency N:N is a NarrowDependencyModified the description of NarrowDependency into 3 different cases with detaild explaination, clearer than the 2 cases explaination before
@zzl0Fisrt four chaptersLots of typos,such as "groupByKey has generated the 3 following RDDs",should be 2. Check pull requestAll fixed
@左手牵右手TELCache and Broadcast chapterLots of typosAll fixed
@cloud-fanJobLogicalPlanSome arrows in the Cogroup() diagram should be colored redAll fixed
@CrazyJvmShuffle detailsStarting from Spark 1.1, the default value for spark.shuffle.file.buffer.kb is 32k, not 100kAll fixed

Special thanks to @明风Andy for his great support.

Special thanks to the rockers (including researchers, developers and users) who participate in the design, implementation and discussion of big data systems.

About

Notes talking about the design and implementation of Apache Spark

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

, '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

Spark Internals

Spark Version: 1.0.2
Doc Version: 1.0.2.0

Authors

Weibo/Twitter IDNameContributions
@JerryLeadLijie XuAuthor of the original Chinese version, and English version update
@juhanlolHan JUEnglish version and update (Chapter 0, 1, 3, 4, and 7)
@invkrhHao RenEnglish version and update (Chapter 2, 5, and 6)

Introduction

This series discuss the design and implementation of Apache Spark, with focuses on its design principles, execution mechanisms, system architecture and performance optimization. In addition, there's some comparisons with Hadoop MapReduce in terms of design and implementation. I'm reluctant to call this document a "code walkthrough", because the goal is not to analyze each piece of code in the project, but to understand the whole system in a systematic way (through analyzing the execution procedure of a Spark job, from its creation to completion).

There're many ways to discuss a computer system. Here, We've chosen a problem-driven approach. Firstly one concret problem is introduced, then it gets analyzed step by step. We'll start from a typical Spark example job and then discuss all the related important system modules. I believe that this approach is better than diving into each module right from the beginning.

The target audiences of this series are geeks who want to have a deeper understanding of Apache Spark as well as other distributed computing frameworks.

I'll try my best to keep this documentation up to date with Spark since it's a fast evolving project with an active community. The documentation's main version is in sync with Spark's version. The additional number at the end represents the documentation's update version.

For more academic oriented discussion, please check out Matei's PHD thesis and other related papers. You can also have a look at my blog (in Chinese) blog.

I haven't been writing such complete documentation for a while. Last time it was about three years ago when I'm studying Andrew Ng's ML course. I was really motivated at that time! This time I've spent 20+ days on this document, from the summer break till now (August 2014). Most of the time is spent on debugging, drawing diagrams and thinking how to put my ideas in the right way. I hope you find this series helpful.

Contents

We start from the creation of a Spark job, and then discuss its execution. Finally, we dive into some related system modules and features.

  1. Overview Overview of Apache Spark
  2. Job logical plan Logical plan of a job (data dependency graph)
  3. Job physical plan Physical plan
  4. Shuffle details Shuffle process
  5. Architecture Coordination of system modules in job execution
  6. Cache and Checkpoint Cache and Checkpoint
  7. Broadcast Broadcast feature
  8. Job Scheduling TODO
  9. Fault-tolerance TODO

Chinses Verison is at markdown/.

The documentation is written in markdown. The pdf version is also available here.

If you're under Max OS X, I recommand MacDown with a github theme for reading.

Gitbook (Chinese version)

Thanks @Yourtion for creating the gitbook version.

Online reading http://spark-internals.books.yourtion.com/

Downloads

Examples

I've created some examples to debug the system during the writing, they are avaible under SparkLearning/src/internals.

Acknowledgement

I appreciate the help from the following in providing solutions and ideas for some detailed issues:

  • @Andrew-Xia Participated in the discussion of BlockManager's implemetation's impact on broadcast(rdd).

  • @CrazyJVM Participated in the discussion of BlockManager's implementation.

  • @王联辉 Participated in the discussion of BlockManager's implementation.

Thanks to the following for complementing the document:

Weibo IDChapterContentRevision status
@OopsOutOfMemoryOverviewRelation between workers and executors and Summary on Spark Executor Driver's Resouce Management (in Chinese)There's not yet a conclusion on this subject since its implementation is still changing, a link to the blog is added

Thanks to the following for finding errors:

Weibo IdChapterError/IssueRevision status
@JoshuawangzjOverviewWhen multiple applications are running, multiple Backend process will be createdCorrected, but need to be confirmed. No idea on how to control the number of Backend processes
@_cs_cmOverviewLatest groupByKey() has removed the mapValues() operation, there's no MapValuesRDD generatedFixed groupByKey() related diagrams and text
@染染生起JobLogicalPlanN:N relation in FullDepedency N:N is a NarrowDependencyModified the description of NarrowDependency into 3 different cases with detaild explaination, clearer than the 2 cases explaination before
@zzl0Fisrt four chaptersLots of typos,such as "groupByKey has generated the 3 following RDDs",should be 2. Check pull requestAll fixed
@左手牵右手TELCache and Broadcast chapterLots of typosAll fixed
@cloud-fanJobLogicalPlanSome arrows in the Cogroup() diagram should be colored redAll fixed
@CrazyJvmShuffle detailsStarting from Spark 1.1, the default value for spark.shuffle.file.buffer.kb is 32k, not 100kAll fixed

Special thanks to @明风Andy for his great support.

Special thanks to the rockers (including researchers, developers and users) who participate in the design, implementation and discussion of big data systems.

About

Notes talking about the design and implementation of Apache Spark

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors

, '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

Spark Internals

Spark Version: 1.0.2
Doc Version: 1.0.2.0

Authors

Weibo/Twitter IDNameContributions
@JerryLeadLijie XuAuthor of the original Chinese version, and English version update
@juhanlolHan JUEnglish version and update (Chapter 0, 1, 3, 4, and 7)
@invkrhHao RenEnglish version and update (Chapter 2, 5, and 6)

Introduction

This series discuss the design and implementation of Apache Spark, with focuses on its design principles, execution mechanisms, system architecture and performance optimization. In addition, there's some comparisons with Hadoop MapReduce in terms of design and implementation. I'm reluctant to call this document a "code walkthrough", because the goal is not to analyze each piece of code in the project, but to understand the whole system in a systematic way (through analyzing the execution procedure of a Spark job, from its creation to completion).

There're many ways to discuss a computer system. Here, We've chosen a problem-driven approach. Firstly one concret problem is introduced, then it gets analyzed step by step. We'll start from a typical Spark example job and then discuss all the related important system modules. I believe that this approach is better than diving into each module right from the beginning.

The target audiences of this series are geeks who want to have a deeper understanding of Apache Spark as well as other distributed computing frameworks.

I'll try my best to keep this documentation up to date with Spark since it's a fast evolving project with an active community. The documentation's main version is in sync with Spark's version. The additional number at the end represents the documentation's update version.

For more academic oriented discussion, please check out Matei's PHD thesis and other related papers. You can also have a look at my blog (in Chinese) blog.

I haven't been writing such complete documentation for a while. Last time it was about three years ago when I'm studying Andrew Ng's ML course. I was really motivated at that time! This time I've spent 20+ days on this document, from the summer break till now (August 2014). Most of the time is spent on debugging, drawing diagrams and thinking how to put my ideas in the right way. I hope you find this series helpful.

Contents

We start from the creation of a Spark job, and then discuss its execution. Finally, we dive into some related system modules and features.

  1. Overview Overview of Apache Spark
  2. Job logical plan Logical plan of a job (data dependency graph)
  3. Job physical plan Physical plan
  4. Shuffle details Shuffle process
  5. Architecture Coordination of system modules in job execution
  6. Cache and Checkpoint Cache and Checkpoint
  7. Broadcast Broadcast feature
  8. Job Scheduling TODO
  9. Fault-tolerance TODO

Chinses Verison is at markdown/.

The documentation is written in markdown. The pdf version is also available here.

If you're under Max OS X, I recommand MacDown with a github theme for reading.

Gitbook (Chinese version)

Thanks @Yourtion for creating the gitbook version.

Online reading http://spark-internals.books.yourtion.com/

Downloads

Examples

I've created some examples to debug the system during the writing, they are avaible under SparkLearning/src/internals.

Acknowledgement

I appreciate the help from the following in providing solutions and ideas for some detailed issues:

  • @Andrew-Xia Participated in the discussion of BlockManager's implemetation's impact on broadcast(rdd).

  • @CrazyJVM Participated in the discussion of BlockManager's implementation.

  • @王联辉 Participated in the discussion of BlockManager's implementation.

Thanks to the following for complementing the document:

Weibo IDChapterContentRevision status
@OopsOutOfMemoryOverviewRelation between workers and executors and Summary on Spark Executor Driver's Resouce Management (in Chinese)There's not yet a conclusion on this subject since its implementation is still changing, a link to the blog is added

Thanks to the following for finding errors:

Weibo IdChapterError/IssueRevision status
@JoshuawangzjOverviewWhen multiple applications are running, multiple Backend process will be createdCorrected, but need to be confirmed. No idea on how to control the number of Backend processes
@_cs_cmOverviewLatest groupByKey() has removed the mapValues() operation, there's no MapValuesRDD generatedFixed groupByKey() related diagrams and text
@染染生起JobLogicalPlanN:N relation in FullDepedency N:N is a NarrowDependencyModified the description of NarrowDependency into 3 different cases with detaild explaination, clearer than the 2 cases explaination before
@zzl0Fisrt four chaptersLots of typos,such as "groupByKey has generated the 3 following RDDs",should be 2. Check pull requestAll fixed
@左手牵右手TELCache and Broadcast chapterLots of typosAll fixed
@cloud-fanJobLogicalPlanSome arrows in the Cogroup() diagram should be colored redAll fixed
@CrazyJvmShuffle detailsStarting from Spark 1.1, the default value for spark.shuffle.file.buffer.kb is 32k, not 100kAll fixed

Special thanks to @明风Andy for his great support.

Special thanks to the rockers (including researchers, developers and users) who participate in the design, implementation and discussion of big data systems.

About

Notes talking about the design and implementation of Apache Spark

Resources

Stars

0 stars

Watchers

1 watching

Forks

Releases

Packages

Contributors