
<!DOCTYPE html>
<html lang="en" dir=ZgotmplZ>

<head>
  


<link rel="stylesheet" href="/bootstrap/css/bootstrap.min.css">
<script src="/bootstrap/js/bootstrap.bundle.min.js"></script>
<link rel="stylesheet" type="text/css" href="/font-awesome/css/font-awesome.min.css">
<script src="/js/anchor.min.js"></script>
<script src="/js/flink.js"></script>
<link rel="canonical" href="https://flink.apache.org/2020/02/22/apache-beam-how-beam-runs-on-top-of-flink/">

  <meta charset="UTF-8">
<meta name="viewport" content="width=device-width, initial-scale=1.0">
<meta name="description" content="Note: This blog post is based on the talk &ldquo;Beam on Flink: How Does It Actually Work?&rdquo;.
Apache Flink and Apache Beam are open-source frameworks for parallel, distributed data processing at scale. Unlike Flink, Beam does not come with a full-blown execution engine of its own but plugs into other execution engines, such as Apache Flink, Apache Spark, or Google Cloud Dataflow. In this blog post we discuss the reasons to use Flink together with Beam for your batch and stream processing needs.">
<meta name="theme-color" content="#FFFFFF"><meta property="og:title" content="Apache Beam: How Beam Runs on Top of Flink" />
<meta property="og:description" content="Note: This blog post is based on the talk &ldquo;Beam on Flink: How Does It Actually Work?&rdquo;.
Apache Flink and Apache Beam are open-source frameworks for parallel, distributed data processing at scale. Unlike Flink, Beam does not come with a full-blown execution engine of its own but plugs into other execution engines, such as Apache Flink, Apache Spark, or Google Cloud Dataflow. In this blog post we discuss the reasons to use Flink together with Beam for your batch and stream processing needs." />
<meta property="og:type" content="article" />
<meta property="og:url" content="https://flink.apache.org/2020/02/22/apache-beam-how-beam-runs-on-top-of-flink/" /><meta property="article:section" content="posts" />
<meta property="article:published_time" content="2020-02-22T12:00:00+00:00" />
<meta property="article:modified_time" content="2020-02-22T12:00:00+00:00" />
<title>Apache Beam: How Beam Runs on Top of Flink | Apache Flink</title>
<link rel="manifest" href="/manifest.json">
<link rel="icon" href="/favicon.png" type="image/x-icon">
<link rel="stylesheet" href="/book.min.22eceb4d17baa9cdc0f57345edd6f215a40474022dfee39b63befb5fb3c596b5.css" integrity="sha256-IuzrTRe6qc3A9XNF7dbyFaQEdAIt/uObY777X7PFlrU=">
<script defer src="/en.search.min.2698f0d1b683dae4d6cb071668b310a55ebcf1c48d11410a015a51d90105b53e.js" integrity="sha256-Jpjw0baD2uTWywcWaLMQpV688cSNEUEKAVpR2QEFtT4="></script>
<!--
Made with Book Theme
https://github.com/alex-shpak/hugo-book
-->

  <meta name="generator" content="Hugo 0.124.1">

    
    <script>
      var _paq = window._paq = window._paq || [];
       
       
      _paq.push(['disableCookies']);
       
      _paq.push(["setDomains", ["*.flink.apache.org","*.nightlies.apache.org/flink"]]);
      _paq.push(['trackPageView']);
      _paq.push(['enableLinkTracking']);
      (function() {
        var u="//analytics.apache.org/";
        _paq.push(['setTrackerUrl', u+'matomo.php']);
        _paq.push(['setSiteId', '1']);
        var d=document, g=d.createElement('script'), s=d.getElementsByTagName('script')[0];
        g.async=true; g.src=u+'matomo.js'; s.parentNode.insertBefore(g,s);
      })();
    </script>
    
</head>

<body dir=ZgotmplZ>
  


<header>
  <nav class="navbar navbar-expand-xl">
    <div class="container-fluid">
      <a class="navbar-brand" href="/">
        <img src="/img/logo/png/100/flink_squirrel_100_color.png" alt="Apache Flink" height="47" width="47" class="d-inline-block align-text-middle">
        <span>Apache Flink</span>
      </a>
      <button class="navbar-toggler" type="button" data-bs-toggle="collapse" data-bs-target="#navbarSupportedContent" aria-controls="navbarSupportedContent" aria-expanded="false" aria-label="Toggle navigation">
          <i class="fa fa-bars navbar-toggler-icon"></i>
      </button>
      <div class="collapse navbar-collapse" id="navbarSupportedContent">
        <ul class="navbar-nav">
          





    
      
  
    <li class="nav-item dropdown">
      <a class="nav-link dropdown-toggle" href="#" role="button" data-bs-toggle="dropdown" aria-expanded="false">About</a>
      <ul class="dropdown-menu">
        
          <li>
            
  
    <a class="dropdown-item" href="/what-is-flink/flink-architecture/">Architecture</a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="/what-is-flink/flink-applications/">Applications</a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="/what-is-flink/flink-operations/">Operations</a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="/what-is-flink/use-cases/">Use Cases</a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="/what-is-flink/powered-by/">Powered By</a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="/what-is-flink/roadmap/">Roadmap</a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="/what-is-flink/community/">Community & Project Info</a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="/what-is-flink/security/">Security</a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="/what-is-flink/special-thanks/">Special Thanks</a>
  

          </li>
        
      </ul>
    </li>
  

    
      
  
    <li class="nav-item dropdown">
      <a class="nav-link dropdown-toggle" href="#" role="button" data-bs-toggle="dropdown" aria-expanded="false">Getting Started</a>
      <ul class="dropdown-menu">
        
          <li>
            
  
    <a class="dropdown-item" href="https://nightlies.apache.org/flink/flink-docs-stable/docs/try-flink/local_installation/">With Flink<i class="link fa fa-external-link title" aria-hidden="true"></i>
    </a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="https://nightlies.apache.org/flink/flink-kubernetes-operator-docs-stable/docs/try-flink-kubernetes-operator/quick-start/">With Flink Kubernetes Operator<i class="link fa fa-external-link title" aria-hidden="true"></i>
    </a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="https://nightlies.apache.org/flink/flink-cdc-docs-stable/docs/get-started/introduction/">With Flink CDC<i class="link fa fa-external-link title" aria-hidden="true"></i>
    </a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="https://nightlies.apache.org/flink/flink-ml-docs-stable/docs/try-flink-ml/quick-start/">With Flink ML<i class="link fa fa-external-link title" aria-hidden="true"></i>
    </a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="https://nightlies.apache.org/flink/flink-statefun-docs-stable/getting-started/project-setup.html">With Flink Stateful Functions<i class="link fa fa-external-link title" aria-hidden="true"></i>
    </a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="https://nightlies.apache.org/flink/flink-docs-stable/docs/learn-flink/overview/">Training Course<i class="link fa fa-external-link title" aria-hidden="true"></i>
    </a>
  

          </li>
        
      </ul>
    </li>
  

    
      
  
    <li class="nav-item dropdown">
      <a class="nav-link dropdown-toggle" href="#" role="button" data-bs-toggle="dropdown" aria-expanded="false">Documentation</a>
      <ul class="dropdown-menu">
        
          <li>
            
  
    <a class="dropdown-item" href="https://nightlies.apache.org/flink/flink-docs-stable/">Flink 1.19 (stable)<i class="link fa fa-external-link title" aria-hidden="true"></i>
    </a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="https://nightlies.apache.org/flink/flink-docs-master/">Flink Master (snapshot)<i class="link fa fa-external-link title" aria-hidden="true"></i>
    </a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="https://nightlies.apache.org/flink/flink-kubernetes-operator-docs-stable/">Kubernetes Operator 1.8 (latest)<i class="link fa fa-external-link title" aria-hidden="true"></i>
    </a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="https://nightlies.apache.org/flink/flink-kubernetes-operator-docs-main">Kubernetes Operator Main (snapshot)<i class="link fa fa-external-link title" aria-hidden="true"></i>
    </a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="https://nightlies.apache.org/flink/flink-cdc-docs-stable">CDC 3.0 (stable)<i class="link fa fa-external-link title" aria-hidden="true"></i>
    </a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="https://nightlies.apache.org/flink/flink-cdc-docs-master">CDC Master (snapshot)<i class="link fa fa-external-link title" aria-hidden="true"></i>
    </a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="https://nightlies.apache.org/flink/flink-ml-docs-stable/">ML 2.3 (stable)<i class="link fa fa-external-link title" aria-hidden="true"></i>
    </a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="https://nightlies.apache.org/flink/flink-ml-docs-master">ML Master (snapshot)<i class="link fa fa-external-link title" aria-hidden="true"></i>
    </a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="https://nightlies.apache.org/flink/flink-statefun-docs-stable/">Stateful Functions 3.3 (stable)<i class="link fa fa-external-link title" aria-hidden="true"></i>
    </a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="https://nightlies.apache.org/flink/flink-statefun-docs-master">Stateful Functions Master (snapshot)<i class="link fa fa-external-link title" aria-hidden="true"></i>
    </a>
  

          </li>
        
      </ul>
    </li>
  

    
      
  
    <li class="nav-item dropdown">
      <a class="nav-link dropdown-toggle" href="#" role="button" data-bs-toggle="dropdown" aria-expanded="false">How to Contribute</a>
      <ul class="dropdown-menu">
        
          <li>
            
  
    <a class="dropdown-item" href="/how-to-contribute/overview/">Overview</a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="/how-to-contribute/contribute-code/">Contribute Code</a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="/how-to-contribute/reviewing-prs/">Review Pull Requests</a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="/how-to-contribute/code-style-and-quality-preamble/">Code Style and Quality Guide</a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="/how-to-contribute/contribute-documentation/">Contribute Documentation</a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="/how-to-contribute/documentation-style-guide/">Documentation Style Guide</a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="/how-to-contribute/improve-website/">Contribute to the Website</a>
  

          </li>
        
          <li>
            
  
    <a class="dropdown-item" href="/how-to-contribute/getting-help/">Getting Help</a>
  

          </li>
        
      </ul>
    </li>
  

    


    
      
  
    <li class="nav-item">
      
  
    <a class="nav-link" href="/posts/">Flink Blog</a>
  

    </li>
  

    
      
  
    <li class="nav-item">
      
  
    <a class="nav-link" href="/downloads/">Downloads</a>
  

    </li>
  

    


    









        </ul>
        <div class="book-search">
          <div class="book-search-spinner hidden">
            <i class="fa fa-refresh fa-spin"></i>
          </div>
          <form class="search-bar d-flex" onsubmit="return false;"su>
            <input type="text" id="book-search-input" placeholder="Search" aria-label="Search" maxlength="64" data-hotkeys="s/">
            <i class="fa fa-search search"></i>
            <i class="fa fa-circle-o-notch fa-spin spinner"></i>
          </form>
          <div class="book-search-spinner hidden"></div>
          <ul id="book-search-results"></ul>
        </div>
      </div>
    </div>
  </nav>
  <div class="navbar-clearfix"></div>
</header>
 
  
      <main class="flex">
        <section class="container book-page">
          
<article class="markdown">
    <h1>
        <a href="/2020/02/22/apache-beam-how-beam-runs-on-top-of-flink/">Apache Beam: How Beam Runs on Top of Flink</a>
    </h1>
    


  February 22, 2020 -



  Maximilian Michels

  <a href="https://twitter.com/stadtlegende">(@stadtlegende)</a>
  

  Markos Sfikas

  <a href="https://twitter.com/MarkSfik">(@MarkSfik)</a>
  



    <p><p>Note: This blog post is based on the talk <a href="https://www.youtube.com/watch?v=hxHGLrshnCY">&ldquo;Beam on Flink: How Does It Actually Work?&rdquo;</a>.</p>
<p><a href="https://flink.apache.org/">Apache Flink</a> and <a href="https://beam.apache.org/">Apache Beam</a> are open-source frameworks for parallel, distributed data processing at scale. Unlike Flink, Beam does not come with a full-blown execution engine of its own but plugs into other execution engines, such as Apache Flink, Apache Spark, or Google Cloud Dataflow. In this blog post we discuss the reasons to use Flink together with Beam for your batch and stream processing needs. We also take a closer look at how Beam works with Flink to provide an idea of the technical aspects of running Beam pipelines with Flink. We hope you find some useful information on how and why the two frameworks can be utilized in combination. For more information, you can refer to the corresponding <a href="https://beam.apache.org/documentation/runners/flink/">documentation</a> on the Beam website or contact the community through the <a href="https://beam.apache.org/community/contact-us/">Beam mailing list</a>.</p>
<h1 id="what-is-apache-beam">
  What is Apache Beam
  <a class="anchor" href="#what-is-apache-beam">#</a>
</h1>
<p><a href="https://beam.apache.org/">Apache Beam</a> is an open-source, unified model for defining batch and streaming data-parallel processing pipelines. It is unified in the sense that you use a single API, in contrast to using a separate API for batch and streaming like it is the case in Flink. Beam was originally developed by Google which released it in 2014 as the Cloud Dataflow SDK. In 2016, it was donated to <a href="https://www.apache.org/">the Apache Software Foundation</a> with the name of Beam. It has been developed by the open-source community ever since. With Apache Beam, developers can write data processing jobs, also known as pipelines, in multiple languages, e.g. Java, Python, Go, SQL. A pipeline is then executed by one of Beam’s Runners. A Runner is responsible for translating Beam pipelines such that they can run on an execution engine. Every supported execution engine has a Runner. The following Runners are available: Apache Flink, Apache Spark, Apache Samza, Hazelcast Jet, Google Cloud Dataflow, and others.</p>
<p>The execution model, as well as the API of Apache Beam, are similar to Flink&rsquo;s. Both frameworks are inspired by the <a href="https://static.googleusercontent.com/media/research.google.com/en//archive/mapreduce-osdi04.pdf">MapReduce</a>, <a href="https://static.googleusercontent.com/media/research.google.com/en//pubs/archive/41378.pdf">MillWheel</a>, and <a href="https://research.google/pubs/pub43864/">Dataflow</a> papers. Like Flink, Beam is designed for parallel, distributed data processing. Both have similar transformations, support for windowing, event/processing time, watermarks, timers, triggers, and much more. However, Beam not being a full runtime focuses on providing the framework for building portable, multi-language batch and stream processing pipelines such that they can be run across several execution engines. The idea is that you write your pipeline once and feed it with either batch or streaming data. When you run it, you just pick one of the supported backends to execute. A large integration test suite in Beam called &ldquo;ValidatesRunner&rdquo; ensures that the results will be the same, regardless of which backend you choose for the execution.</p>
<p>One of the most exciting developments in the Beam technology is the framework’s support for multiple programming languages including Java, Python, Go, Scala and SQL. Essentially, developers can write their applications in a programming language of their choice. Beam, with the help of the Runners, translates the program to one of the execution engines, as shown in the diagram below.</p>
<center>
<img src="/img/blog/2020-02-22-beam-on-flink/flink-runner-beam-beam-vision.png" width="600px" alt="The vision of Apache Beam"/>
</center>
<h1 id="reasons-to-use-beam-with-flink">
  Reasons to use Beam with Flink
  <a class="anchor" href="#reasons-to-use-beam-with-flink">#</a>
</h1>
<p>Why would you want to use Beam with Flink instead of directly using Flink? Ultimately, Beam and Flink complement each other and provide additional value to the user. The main reasons for using Beam with Flink are the following:</p>
<ul>
<li>Beam provides a unified API for both batch and streaming scenarios.</li>
<li>Beam comes with native support for different programming languages, like Python or Go with all their libraries like Numpy, Pandas, Tensorflow, or TFX.</li>
<li>You get the power of Apache Flink like its exactly-once semantics, strong memory management and robustness.</li>
<li>Beam programs run on your existing Flink infrastructure or infrastructure for other supported Runners, like Spark or Google Cloud Dataflow.</li>
<li>You get additional features like side inputs and cross-language pipelines that are not supported natively in Flink but only supported when using Beam with Flink.</li>
</ul>
<h1 id="the-flink-runner-in-beam">
  The Flink Runner in Beam
  <a class="anchor" href="#the-flink-runner-in-beam">#</a>
</h1>
<p>The Flink Runner in Beam translates Beam pipelines into Flink jobs. The translation can be parameterized using Beam&rsquo;s pipeline options which are parameters for settings like configuring the job name, parallelism, checkpointing, or metrics reporting.</p>
<p>If you are familiar with a DataSet or a DataStream, you will have no problems understanding what a PCollection is. PCollection stands for parallel collection in Beam and is exactly what DataSet/DataStream would be in Flink. Due to Beam&rsquo;s unified API we only have one type of results of transformation: PCollection.</p>
<p>Beam pipelines are composed of transforms. Transforms are like operators in Flink and come in two flavors: primitive and composite transforms. The beauty of all this is that Beam only comes with a small set of primitive transforms which are:</p>
<ul>
<li><code>Source</code> (for loading data)</li>
<li><code>ParDo</code> (think of a flat map operator on steroids)</li>
<li><code>GroupByKey</code> (think of keyBy() in Flink)</li>
<li><code>AssignWindows</code> (windows can be assigned at any point in time in Beam)</li>
<li><code>Flatten</code> (like a union() operation in Flink)</li>
</ul>
<p>Composite transforms are built by combining the above primitive transforms. For example, <code>Combine = GroupByKey + ParDo</code>.</p>
<h1 id="flink-runner-internals">
  Flink Runner Internals
  <a class="anchor" href="#flink-runner-internals">#</a>
</h1>
<p>Although using the Flink Runner in Beam has no prerequisite to understanding its internals, we provide more details of how the Flink runner works in Beam to share knowledge of how the two frameworks can integrate and work together to provide state-of-the-art streaming data pipelines.</p>
<p>The Flink Runner has two translation paths. Depending on whether we execute in batch or streaming mode, the Runner either translates into Flink&rsquo;s DataSet or into Flink&rsquo;s DataStream API. Since multi-language support has been added to Beam, another two translation paths have been added. To summarize the four modes:</p>
<ol>
<li><strong>The Classic Flink Runner for batch jobs:</strong> Executes batch Java pipelines</li>
<li><strong>The Classic Flink Runner for streaming jobs:</strong> Executes streaming Java pipelines</li>
<li><strong>The Portable Flink Runner for batch jobs:</strong> Executes Java as well as Python, Go and other supported SDK pipelines for batch scenarios</li>
<li><strong>The Portable Flink Runner for streaming jobs:</strong> Executes Java as well as Python, Go and other supported SDK pipelines for streaming scenarios</li>
</ol>
<center>
<img src="/img/blog/2020-02-22-beam-on-flink/flink-runner-beam-runner-translation-paths.png" width="300px" alt="The 4 translation paths in the Beam's Flink Runner"/>
</center>
<h2 id="the-classic-flink-runner-in-beam">
  The “Classic” Flink Runner in Beam
  <a class="anchor" href="#the-classic-flink-runner-in-beam">#</a>
</h2>
<p>The classic Flink Runner was the initial version of the Runner, hence the &ldquo;classic&rdquo; name. Beam pipelines are represented as a graph in Java which is composed of the aforementioned composite and primitive transforms. Beam provides translators which traverse the graph in topological order. Topological order means that we start from all the sources first as we iterate through the graph. Presented with a transform from the graph, the Flink Runner generates the API calls as you would normally when writing a Flink job.</p>
<center>
<img src="/img/blog/2020-02-22-beam-on-flink/classic-flink-runner-beam.png" width="600px" alt="The Classic Flink Runner in Beam"/>
</center>
<p>While Beam and Flink share very similar concepts, there are enough differences between the two frameworks that make Beam pipelines impossible to be translated 1:1 into a Flink program. In the following sections, we will present the key differences:</p>
<h3 id="serializers-vs-coders">
  Serializers vs Coders
  <a class="anchor" href="#serializers-vs-coders">#</a>
</h3>
<p>When data is transferred over the wire in Flink, it has to be turned into bytes. This is done with the help of serializers. Flink has a type system to instantiate the correct coder for a given type, e.g. <code>StringTypeSerializer</code> for a String. Apache Beam also has its own type system which is similar to Flink&rsquo;s but uses slightly different interfaces. Serializers are called Coders in Beam. In order to make a Beam Coder run in Flink, we have to make the two serializer types compatible. This is done by creating a special Flink type information that looks like the one in Flink but calls the appropriate Beam coder. That way, we can use Beam&rsquo;s coders although we are executing the Beam job with Flink. Flink operators expect a TypeInformation, e.g. <code>StringTypeInformation</code>, for which we use a <code>CoderTypeInformation</code> in Beam. The type information returns the serializer for which we return a <code>CoderTypeSerializer</code>, which calls the underlying Beam Coder.</p>
<center>
<img src="/img/blog/2020-02-22-beam-on-flink/flink-runner-beam-serializers-coders.png" width="300px" alt="Serializers vs Coders"/>
</center>
<h3 id="read">
  Read
  <a class="anchor" href="#read">#</a>
</h3>
<p>The <code>Read</code> transform provides a way to read data into your pipeline in Beam. The Read transform is supported by two wrappers in Beam, the <code>SourceInputFormat</code> for batch processing and the <code>UnboundedSourceWrapper</code> for stream processing.</p>
<h3 id="pardo">
  ParDo
  <a class="anchor" href="#pardo">#</a>
</h3>
<p><code>ParDo</code> is the swiss army knife of Beam and can be compared to a <code>RichFlatMapFunction</code> in Flink with additional features such as <code>SideInputs</code>, <code>SideOutputs</code>, State and Timers. <code>ParDo</code> is essentially translated by the Flink runner using the <code>FlinkDoFnFunction</code> for batch processing or the <code>FlinkStatefulDoFnFunction</code>, while for streaming scenarios the translation is executed with the <code>DoFnOperator</code> that takes care of checkpointing and buffering of data during checkpoints, watermark emissions and maintenance of state and timers. This is all executed by Beam’s interface, called the <code>DoFnRunner</code>, that encapsulates Beam-specific execution logic, like retrieving state, executing state and timers, or reporting metrics.</p>
<h3 id="side-inputs">
  Side Inputs
  <a class="anchor" href="#side-inputs">#</a>
</h3>
<p>In addition to the main input, ParDo transforms can have a number of side inputs. A side input can be a static set of data that you want to have available at all parallel instances. However, it is more flexible than that. You can have keyed and even windowed side input which updates based on the window size. This is a very powerful concept which does not exist in Flink but is added on top of Flink using Beam.</p>
<h3 id="assignwindows">
  AssignWindows
  <a class="anchor" href="#assignwindows">#</a>
</h3>
<p>In Flink, windows are assigned by the <code>WindowOperator</code> when you use the <code>window()</code> in the API. In Beam, windows can be assigned at any point in time. Any element is implicitly part of a window. If no window is assigned explicitly, the element is part of the <code>GlobalWindow</code>. Window information is stored for each element in a wrapper called <code>WindowedValue</code>. The window information is only used once we issue a <code>GroupByKey</code>.</p>
<h3 id="groupbykey">
  GroupByKey
  <a class="anchor" href="#groupbykey">#</a>
</h3>
<p>Most of the time it is useful to partition the data by a key. In Flink, this is done via the <code>keyBy()</code> API call. In Beam the <code>GroupByKey</code> transform can only be applied if the input is of the form <code>KV&lt;Key, Value&gt;</code>. Unlike Flink where the key can even be nested inside the data, Beam enforces the key to always be explicit. The <code>GroupByKey</code> transform then groups the data by key and by window which is similar to what <code>keyBy(..).window(..)</code> would give us in Flink. Beam has its own set of libraries to do that because Beam has its own set of window functions and triggers. Essentially, GroupByKey is very similar to what the WindowOperator does in Flink.</p>
<h3 id="flatten">
  Flatten
  <a class="anchor" href="#flatten">#</a>
</h3>
<p>The Flatten operator takes multiple DataSet/DataStreams, called P[arallel]Collections in Beam, and combines them into one collection. This is equivalent to Flink&rsquo;s <code>union()</code> operation.</p>
<h2 id="the-portable-flink-runner-in-beam">
  The “Portable” Flink Runner in Beam
  <a class="anchor" href="#the-portable-flink-runner-in-beam">#</a>
</h2>
<p>The portable Flink Runner in Beam is the evolution of the classic Runner. Classic Runners are tied to the JVM ecosystem, but the Beam community wanted to move past this and also execute Python, Go and other languages. This adds another dimension to Beam in terms of portability because, like previously mentioned, Beam already had portability across execution engines. It was necessary to change the translation logic of the Runner to be able to support language portability.</p>
<p>There are two important building blocks for portable Runners:</p>
<ol>
<li>A common pipeline format across all the languages: The Runner API</li>
<li>A common interface during execution for the communication between the Runner and the code written in any language: The Fn API</li>
</ol>
<p>The Runner API provides a universal representation of the pipeline as Protobuf which contains the transforms, types, and user code. Protobuf was chosen as the format because every language has libraries available for it. Similarly, for the execution part, Beam introduced the Fn API interface to handle the communication between the Runner/execution engine and the user code that may be written in a different language and executes in a different process. Fn API is pronounced &ldquo;fun API&rdquo;, you may guess why.</p>
<center>
<img src="/img/blog/2020-02-22-beam-on-flink/flink-runner-beam-language-portability.png" width="600px" alt="Language Portability in Apache Beam"/>
</center>
<h2 id="how-are-beam-programs-translated-in-language-portability">
  How Are Beam Programs Translated In Language Portability?
  <a class="anchor" href="#how-are-beam-programs-translated-in-language-portability">#</a>
</h2>
<p>Users write their Beam pipelines in one language, but they may get executed in an environment based on a completely different language. How does that work? To explain that, let&rsquo;s follow the lifecycle of a pipeline. Let&rsquo;s suppose we use the Python SDK to write the pipeline. Before submitting the pipeline via the Job API to Beam&rsquo;s JobServer, Beam would convert it to the Runner API, the language-agnostic format we described before. The JobServer is also a Beam component that handles the staging of the required dependencies during execution. The JobServer will then kick-off the translation which is similar to the classic Runner. However, an important change is the so-called <code>ExecutableStage</code> transform. It is essentially a ParDo transform that we already know but designed for holding language-dependent code. Beam tries to combine as many of these transforms into one &ldquo;executable stage&rdquo;. The result again is a Flink program which is then sent to the Flink cluster and executed there. The major difference compared to the classic Runner is that during execution we will start <em>environments</em> to execute the aforementioned <em>ExecutableStages</em>. The following environments are available:</p>
<ul>
<li>Docker-based (the default)</li>
<li>Process-based (a simple process is started)</li>
<li>Externally-provided (K8s or other schedulers)</li>
<li>Embedded (intended for testing and only works with Java)</li>
</ul>
<p>Environments hold the <em>SDK Harness</em> which is the code that handles the execution and the communication with the Runner over the Fn API. For example, when Flink executes Python code, it sends the data to the Python environment containing the Python SDK Harness. Sending data to an external process involves a minor overhead which we have measured to be 5-10% slower than the classic Java pipelines. However, Beam uses a fusion of transforms to execute as many transforms as possible in the same environment which share the same input or output. That&rsquo;s why in real-world scenarios the overhead could be much lower.</p>
<center>
<img src="/img/blog/2020-02-22-beam-on-flink/flink-runner-beam-language-portability-architecture.png" width="600px" alt="Language Portability Architecture in beam"/>
</center>
<p>Environments can be present for many languages. This opens up an entirely new type of pipelines: cross-language pipelines. In cross-language pipelines we can combine transforms of two or more languages, e.g. a machine learning pipeline with the feature generation written in Java and the learning written in Python. All this can be run on top of Flink.</p>
<h2 id="conclusion">
  Conclusion
  <a class="anchor" href="#conclusion">#</a>
</h2>
<p>Using Apache Beam with Apache Flink combines  (a.) the power of Flink with (b.) the flexibility of Beam. All it takes to run Beam is a Flink cluster, which you may already have. Apache Beam&rsquo;s fully-fledged Python API is probably the most compelling argument for using Beam with Flink, but the unified API which allows to &ldquo;write-once&rdquo; and &ldquo;execute-anywhere&rdquo; is also very appealing to Beam users. On top of this, features like side inputs and a rich connector ecosystem are also reasons why people like Beam.</p>
<p>With the introduction of schemas, a new format for handling type information, Beam is heading in a similar direction as Flink with its type system which is essential for the Table API or SQL. Speaking of, the next Flink release will include a Python version of the Table API which is based on the language portability of Beam. Looking ahead, the Beam community plans to extend the support for interactive programs like notebooks. TFX, which is built with Beam, is a very powerful way to solve many problems around training and validating machine learning models.</p>
<p>For many years, Beam and Flink have inspired and learned from each other. With the Python support being based on Beam in Flink, they only seem to come closer to each other. That&rsquo;s all the better for the community, and also users have more options and functionality to choose from.</p>
</p>
</article>

          



  
    
    <div class="edit-this-page">
      <p>
        <a href="https://cwiki.apache.org/confluence/display/FLINK/Flink+Translation+Specifications">Want to contribute translation?</a>
      </p>
      <p>
        <a href="//github.com/apache/flink-web/edit/asf-site/docs/content/posts/2020-02-22-apache-beam-how-beam-runs-on-top-of-flink.md">
          Edit This Page<i class="fa fa-edit fa-fw"></i> 
        </a>
      </p>
    </div>

        </section>
        
          <aside class="book-toc">
            


<nav id="TableOfContents"><h3>On This Page <a href="javascript:void(0)" class="toc" onclick="collapseToc()"><i class="fa fa-times" aria-hidden="true"></i></a></h3>
  <ul>
    <li><a href="#what-is-apache-beam">What is Apache Beam</a></li>
    <li><a href="#reasons-to-use-beam-with-flink">Reasons to use Beam with Flink</a></li>
    <li><a href="#the-flink-runner-in-beam">The Flink Runner in Beam</a></li>
    <li><a href="#flink-runner-internals">Flink Runner Internals</a>
      <ul>
        <li><a href="#the-classic-flink-runner-in-beam">The “Classic” Flink Runner in Beam</a>
          <ul>
            <li><a href="#serializers-vs-coders">Serializers vs Coders</a></li>
            <li><a href="#read">Read</a></li>
            <li><a href="#pardo">ParDo</a></li>
            <li><a href="#side-inputs">Side Inputs</a></li>
            <li><a href="#assignwindows">AssignWindows</a></li>
            <li><a href="#groupbykey">GroupByKey</a></li>
            <li><a href="#flatten">Flatten</a></li>
          </ul>
        </li>
        <li><a href="#the-portable-flink-runner-in-beam">The “Portable” Flink Runner in Beam</a></li>
        <li><a href="#how-are-beam-programs-translated-in-language-portability">How Are Beam Programs Translated In Language Portability?</a></li>
        <li><a href="#conclusion">Conclusion</a></li>
      </ul>
    </li>
  </ul>
</nav>


          </aside>
          <aside class="expand-toc hidden">
            <a class="toc" onclick="expandToc()" href="javascript:void(0)">
              <i class="fa fa-bars" aria-hidden="true"></i>
            </a>
          </aside>
        
      </main>

      <footer>
        


<div class="separator"></div>
<div class="panels">
  <div class="wrapper">
      <div class="panel">
        <ul>
          <li>
            <a href="https://flink-packages.org/">flink-packages.org</a>
          </li>
          <li>
            <a href="https://www.apache.org/">Apache Software Foundation</a>
          </li>
          <li>
            <a href="https://www.apache.org/licenses/">License</a>
          </li>
          
          
          
            
          
            
          
          

          
            
              
            
          
            
              
                <li>
                  <a  href="/zh/">
                    <i class="fa fa-globe" aria-hidden="true"></i>&nbsp;中文版
                  </a>
                </li>
              
            
          
       </ul>
      </div>
      <div class="panel">
        <ul>
          <li>
            <a href="/what-is-flink/security">Security</a-->
          </li>
          <li>
            <a href="https://www.apache.org/foundation/sponsorship.html">Donate</a>
          </li>
          <li>
            <a href="https://www.apache.org/foundation/thanks.html">Thanks</a>
          </li>
       </ul>
      </div>
      <div class="panel icons">
        <div>
          <a href="/posts">
            <div class="icon flink-blog-icon"></div>
            <span>Flink blog</span>
          </a>
        </div>
        <div>
          <a href="https://github.com/apache/flink">
            <div class="icon flink-github-icon"></div>
            <span>Github</span>
          </a>
        </div>
        <div>
          <a href="https://twitter.com/apacheflink">
            <div class="icon flink-twitter-icon"></div>
            <span>Twitter</span>
          </a>
        </div>
      </div>
  </div>
</div>

<hr/>

<div class="container disclaimer">
  <p>The contents of this website are © 2024 Apache Software Foundation under the terms of the Apache License v2. Apache Flink, Flink, and the Flink logo are either registered trademarks or trademarks of The Apache Software Foundation in the United States and other countries.</p>
</div>



      </footer>
    
  </body>
</html>






