{"id":601885,"date":"2019-11-13T09:41:06","date_gmt":"2019-11-13T08:41:06","guid":{"rendered":"https:\/\/www.devoteam.com\/expert-view\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\/"},"modified":"2019-11-13T09:41:06","modified_gmt":"2019-11-13T08:41:06","slug":"querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam","status":"publish","type":"expert-view","link":"https:\/\/devoteam.info\/en-nl\/expert-view\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\/","title":{"rendered":"Querying JDBC database in parallel with Google Dataflow (Apache Beam)"},"content":{"rendered":"<p><span style=\"font-weight: 400\">\u00a0<\/span><\/p>\n<div class=\"block-fullwidth-dark\">\n<p style=\"text-align: left\"><strong>NOTE: The Java example code for this technical <\/strong><\/p>\n<p style=\"text-align: left\"><strong>blog can be found in this GitHub repo:<\/strong><\/p>\n<div class=\"row-more clearfix\" style=\"text-align: left\"><a class=\"bt-more\" href=\"https:\/\/github.com\/HocLengChung\/Apache-Beam-Jdbc-Parallel-Read\" target=\"_blank\" rel=\"noopener noreferrer\">GitHub Repo<\/a><\/div>\n<\/div>\n<p><span style=\"font-weight: 400\">Consider the following situation: You want to use a single query to query a JDBC compatible database like <a href=\"https:\/\/nl.devoteam.com\/partner\/google-cloud\/\">Google Cloud<\/a> SQL (MySQL) that contains millions of rows. You may want to do this when migrating legacy database data to BigQuery. At some point, the number of rows is so large, that a single query yields too many results for your virtual machines in Google Dataflow to handle. And the pipeline will hang and will not complete at all due to an out of memory problem.<\/span><\/p>\n<h2><span style=\"font-weight: 400\">The Challenge<\/span><\/h2>\n<p><span style=\"font-weight: 400\">The Java version of Apache Beam has the built-in function JdbcIO.read() I\/O Transform that can read and write to a JDBC source. JdbcIO can read the source using a single query. However, if your query returns millions of rows, it will become slow and not complete due to memory issues in your machine. In the Apache Sparks JDBC connector, there is a query partition function that divides a query into partitions before executing it. But, running Spark means that you will need to provision virtual machines on Dataproc before you can run the pipeline. With Apache Beam you can run the pipeline directly using Google Dataflow and any provisioning of machines is done when you specify the pipeline parameters. It is a serverless, on-demand solution. Would it be possible to do something like this in Apache Beam?\u00a0<\/span><\/p>\n<h2><span style=\"font-weight: 400\">Building a partitioned JDBC query pipeline (Java Apache Beam).<\/span><\/h2>\n<p><span style=\"font-weight: 400\">Apache Beams JdbcIO.readAll() Transform can query a source in parallel, given a PCollection of query strings. In order to query a table in parallel, we need to construct queries that query ranges of a table. Consider for example a MySQL table with an auto-increment column \u2018index_id\u2019. See also the SQL DDL command for the example table:<\/span><\/p>\n<div class=\"block-fullwidth-dark\">\n<pre><span style=\"font-weight: 400\">CREATE TABLE `HelloWorld` (<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0`ID` varchar(255) DEFAULT NULL,<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0`Job Title` varchar(255) DEFAULT NULL,<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0`Email Address` varchar(255) DEFAULT NULL,<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0`FirstName LastName` varchar(255) DEFAULT NULL,<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0`Address` varchar(255) DEFAULT NULL,<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0`integer_column` varchar(255) DEFAULT NULL,<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0`string_column` varchar(255) DEFAULT NULL,<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0`index_id` int(11) NOT NULL AUTO_INCREMENT,<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0PRIMARY KEY (`index_id`)<\/span>\n\n<span style=\"font-weight: 400\">)<\/span><\/pre>\n<\/div>\n<p><span style=\"font-weight: 400\">I have populated this table with 16 million rows of data inside a Google Cloud SQL instance.\u00a0<\/span><\/p>\n<h2><span style=\"font-weight: 400\">General step by step implementing partitioned parallel JdbcIO read<\/span><\/h2>\n<p><strong>1) Add dependencies<\/strong><\/p>\n<p><span style=\"font-weight: 400\">A JDBC connection pool will make concurrent JDBC connections more efficient. I use c3p0 for this purpose:<\/span><\/p>\n<div class=\"block-fullwidth-dark\">\n<pre><span style=\"font-weight: 400\">\u00a0\u00a0\u00a0&lt;dependency&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0&lt;groupId&gt;com.mchange&lt;\/groupId&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0&lt;artifactId&gt;c3p0&lt;\/artifactId&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0&lt;version&gt;0.9.5.4&lt;\/version&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0&lt;\/dependency&gt;<\/span><\/pre>\n<\/div>\n<p><span style=\"font-weight: 400\">We need to add JDBC dependencies for the Java Maven project in order to run the example.\u00a0<\/span><\/p>\n<div class=\"block-fullwidth-dark\">\n<pre><span style=\"font-weight: 400\">\u00a0&lt;dependency&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0&lt;groupId&gt;org.apache.beam&lt;\/groupId&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0&lt;artifactId&gt;beam-sdks-java-io-jdbc&lt;\/artifactId&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0&lt;version&gt;${beam.version}&lt;\/version&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0&lt;\/dependency&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0&lt;dependency&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0&lt;groupId&gt;mysql&lt;\/groupId&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0&lt;artifactId&gt;mysql-connector-java&lt;\/artifactId&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0&lt;version&gt;8.0.16&lt;\/version&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0&lt;\/dependency&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0&lt;dependency&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0&lt;groupId&gt;com.google.cloud.sql&lt;\/groupId&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0&lt;artifactId&gt;mysql-socket-factory-connector-j-8&lt;\/artifactId&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0&lt;version&gt;1.0.15&lt;\/version&gt;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0&lt;\/dependency&gt;<\/span><\/pre>\n<\/div>\n<p><strong>2) Configure connection pool of the added dependency mentioned in 1)<\/strong><\/p>\n<div class=\"block-fullwidth-dark\">\n<pre><span style=\"font-weight: 400\">ComboPooledDataSource dataSource = new ComboPooledDataSource();<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0dataSource.setDriverClass(\"com.mysql.cj.jdbc.Driver\");<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0dataSource.setJdbcUrl(\"jdbc:mysql:\/\/google\/&lt;DATABASE_NAME&gt;?cloudSqlInstance=\n&lt;INSTANCE_CONNECTION_NAME&gt;\" +<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\"&amp;socketFactory=com.google.cloud.sql.mysql.SocketFactory&amp;useSSL=false\" +<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\"&amp;user=&lt;MYSQL_USER_NAME&gt;&amp;password=&lt;MYSQL_USER_PASSWORD&gt;\");<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0dataSource.setMaxPoolSize(10);<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0dataSource.setInitialPoolSize(6);<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0JdbcIO.DataSourceConfiguration config<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0= JdbcIO.DataSourceConfiguration.create(dataSource);<\/span><\/pre>\n<\/div>\n<p><span style=\"font-weight: 400\"><strong>3) We can now build the first part of the pipeline<\/strong><\/span><\/p>\n<p><span style=\"font-weight: 400\">Consider the \u201cRead from Cloud SQL MySQL: HelloWorld&#8221; step of the pipeline below. I use the index_id field for pagination. I divide the table into partitioned ranges of chunks of 1000 elements. I use the query:<\/span><\/p>\n<div class=\"block-fullwidth-dark\">\n<pre><span style=\"font-weight: 400\">\u201cSELECT MAX(`index_id`) from HelloWorld\u201d<\/span><\/pre>\n<\/div>\n<p><span style=\"font-weight: 400\">To get the number of rows from the table. This method is faster than using a SELECT COUNT(*) query. Note that this solution is for demonstration purposes only. The MySQL database is using the InnoDB engine which is slow when doing a SELECT COUNT(*) query with a large number of rows. If your table has deleted rows in the past, selecting the maximum value from index_id would not be accurate. In this case, you could store the number of rows of a table in a separate metadata table.<\/span><\/p>\n<p><span style=\"font-weight: 400\">In the \u201cDistribute\u201d step of the pipeline, I convert the total number of rows into Strings of range indices separated by a comma to an Apache Beam Key Value class instance with String as key and Integer as value (KV&lt;String,Integer&gt;). The key value of the KV are the range indices that have been defined separated with a comma and the value of the KV is just a fixed integer 1.\u00a0<\/span><\/p>\n<p><span style=\"font-weight: 400\">In the \u201cBreak Fusion\u201d step I invoke a GroupByKey Transform. This will group the KVs of the earlier step by key. This step is needed in order to \u201cbreak the fusion\u201d. Apparently, if your pipeline contains operations that have a high number of output elements then the elements outputted will stay in the same machine to be processed. You need to redistribute this to all available machines in order to keep the parallel processing capabilities of Apache Beam. See <a href=\"https:\/\/cloud.google.com\/dataflow\/docs\/guides\/deploying-a-pipeline#fusion-optimization\" target=\"_blank\" rel=\"noopener noreferrer\">here<\/a> <\/span><span style=\"font-weight: 400\">for more details on deploying a pipeline.<\/span><\/p>\n<p><strong>The code excerpt of these steps are displayed below:<\/strong><\/p>\n<div class=\"block-fullwidth-dark\">\n<pre><span style=\"font-weight: 400\">PCollection&lt;KV&lt;String,Iterable&lt;Integer&gt;&gt;&gt; ranges =<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0p.apply(String.format(\"Read from Cloud SQL MySQL: %s\",tableName), JdbcIO.\n&lt;String&gt;read()<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0.withDataSourceConfiguration(config)<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0.withQuery(String.format(\"SELECT MAX(`index_id`) from %s\", tableName))<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0.withRowMapper(new JdbcIO.RowMapper&lt;String&gt;() {<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0public String mapRow(ResultSet resultSet) throws Exception {<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0return resultSet.getString(1);<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0}<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0})<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0.withOutputParallelization(false)<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0.withCoder(StringUtf8Coder.of()))<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0.apply(\"Distribute\", ParDo.of(new DoFn&lt;String, KV&lt;String, Integer&gt;&gt;() {<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0@ProcessElement<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0public void processElement(ProcessContext c) {<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0int readChunk = fetchSize;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0int count = Integer.parseInt((String) c.element());<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0int ranges = (int) (count \/ readChunk);<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0for (int i = 0; i &lt; ranges; i++) {<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0int indexFrom = i * readChunk;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0int indexTo = (i + 1) * readChunk;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0String range = String.format(\"%s,%s\",indexFrom, indexTo);<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0KV&lt;String,Integer&gt; kvRange = KV.of(range, 1);<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0c.output(kvRange);<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0}<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0if (count &gt; ranges * readChunk) {<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0int indexFrom = ranges * readChunk;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0int indexTo = ranges * readChunk + count % readChunk;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0String range = String.format(\"%s,%s\",indexFrom, indexTo);<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0KV&lt;String,Integer&gt; kvRange = KV.of(range, 1);<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0c.output(kvRange);<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0}<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0}<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0}))<\/span>\n\n<span style=\"font-weight: 400\"> \u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0.apply(\"Break Fusion\", GroupByKey.create());<\/span><\/pre>\n<\/div>\n<p><span style=\"font-weight: 400\"><strong>4) The last step is to use the created ranges to read JDBC the source in parallel<\/strong>\u00a0<\/span><\/p>\n<p><span style=\"font-weight: 400\">For that I use the query:<\/span><\/p>\n<div class=\"block-fullwidth-dark\">\n<pre><span style=\"font-weight: 400\">select * from &lt;DATABASE_NAME&gt;.HelloWorld where index_id &gt;= ? and index_id &lt; ?<\/span><\/pre>\n<\/div>\n<p><span style=\"font-weight: 400\">and a ParameterSetter. For illustration purposes, I added a RowMapper that maps the selected rows into a readable JSON String using Jackson. I also added a simple map function that does something after reading the data from the database in parallel. The code excerpt for these steps is displayed below:<\/span><\/p>\n<div class=\"block-fullwidth-dark\">\n<pre><span style=\"font-weight: 400\">ranges.apply(String.format(\"Read ALL %s\", tableName), JdbcIO.&lt;KV&lt;String,Iterable&lt;Integer&gt;&gt;,<\/span><\/pre>\n<pre><span style=\"font-weight: 400\">String&gt;readAll()<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0.withDataSourceConfiguration(config)<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0.withFetchSize(fetchSize)<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0.withCoder(StringUtf8Coder.of())<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0.withParameterSetter(new JdbcIO.PreparedStatementSetter&lt;KV&lt;String,Iterable&lt;Integer&gt;&gt;&gt;() {<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0@Override<\/span>\n\n\n\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0public void setParameters(KV&lt;String,Iterable&lt;Integer&gt;&gt; element,<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0PreparedStatement preparedStatement) throws Exception {<\/span>\n\n\n\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0String[] range = element.getKey().split(\",\");<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0preparedStatement.setInt(1, Integer.parseInt(range[0]));<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0preparedStatement.setInt(2, Integer.parseInt(range[1]));<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0}<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0})<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0.withOutputParallelization(false)<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0.withQuery(String.format(\"select * from &lt;DATABASE_NAME&gt;.%s where index_id &gt;= ? <\/span><\/pre>\n<pre><span style=\"font-weight: 400\">              and index_id &lt; ?\",tableName))<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0.withRowMapper((JdbcIO.RowMapper&lt;String&gt;) resultSet -&gt; {<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0ObjectMapper mapper = new ObjectMapper();<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0ArrayNode arrayNode = mapper.createArrayNode();<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0for (int i = 1; i &lt;= resultSet.getMetaData().getColumnCount(); i++) {<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0String columnTypeIntKey =\"\";<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0try {<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0ObjectNode objectNode = mapper.createObjectNode();<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0objectNode.put(\"column_name\",<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0resultSet.getMetaData().getColumnName(i));<\/span>\n\n\n\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0objectNode.put(\"value\",<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0resultSet.getString(i));<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0arrayNode.add(objectNode);<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0} catch (Exception e) {<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0LOG.error(\"problem columnTypeIntKey: \" +\u00a0 columnTypeIntKey);<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0throw e;<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0}<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0}<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0return mapper.writeValueAsString(arrayNode);<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0})<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0)<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0.apply(MapElements.via(<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0new SimpleFunction&lt;String, Integer&gt;() {<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0@Override<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0public Integer apply(String line) {<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0return line.length();<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0}<\/span>\n\n<span style=\"font-weight: 400\">\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0\u00a0}))<\/span>\n\n<span style=\"font-weight: 400\"> \u00a0\u00a0\u00a0\u00a0\u00a0;<\/span><\/pre>\n<\/div>\n<p><span style=\"font-weight: 400\">By now you have a pipeline that reads a JDBC source in parallel. In my example I got a throughput of over 250k elements per second with three n1-standard-8 machines:<\/span><\/p>\n<p><img loading=\"lazy\" decoding=\"async\" class=\"aligncenter wp-image-80199 size-full\" src=\"https:\/\/devoteam.info\/wp-content\/uploads\/2024\/12\/zzzzz.png\" alt=\"JDBC\" width=\"300\" height=\"500\" \/><\/p>\n<h2><span style=\"font-weight: 400\">Conclusion<\/span><\/h2>\n<p><span style=\"font-weight: 400\">In short, this article explained how to read from a JDBC source using JdbcIO.readAll() transform of Apache Beam. The user can use the provided example to implement their own parallel JDBC read pipeline within Apache Beam. If you want to learn more about Apache Beam, click <\/span><a href=\"https:\/\/beam.apache.org\/documentation\/\"><span style=\"font-weight: 400\">here<\/span><\/a><span style=\"font-weight: 400\"> to get all the Apache Beam resources.<\/span><\/p>\n<p><span style=\"font-weight: 400\">If you want to use this solution for your own pipelines please look at the GitHub repo that contains the code: <\/span><a href=\"https:\/\/github.com\/HocLengChung\/Apache-Beam-Jdbc-Parallel-Read\"><span style=\"font-weight: 400\">https:\/\/github.com\/HocLengChung\/Apache-Beam-Jdbc-Parallel-Read<\/span><\/a><\/p>\n<p><span style=\"font-weight: 400\">Just change the parameters of the example to your needs.<\/span><\/p>\n<div class=\"block-fullwidth-dark\">\n<h3>Data Integration At Devoteam<\/h3>\n<p>Data and its value are growing exponentially in the 21st century. Major IT challenges are mostly related to data (and its massive growth). These challenges often include questions such as: How do I get a grip on data? Where is my data stored? Which (advanced) analytics possibilities are there to optimize its use? How reliable or trustworthy is my data? This begs for an integrated data strategy.<\/p>\n<div class=\"row-more clearfix\"><a class=\"bt-more\" href=\"https:\/\/nl.devoteam.com\/playground\/data-driven-intelligence\/\" target=\"_blank\" rel=\"noopener noreferrer\">More Data<\/a><\/div>\n<\/div>\n<div class=\"block-fullwidth-dark\">\n<h1>Join our open &amp; innovative culture<\/h1>\n<p>Open, accessible and ambitious are the keywords of our organizational culture. We encourage making mistakes, we strive to do better every day and love to have fun. Doing work you love with brilliant people is what it&#8217;s all about.<\/p>\n<div class=\"row-more clearfix\"><a class=\"bt-more\" href=\"https:\/\/nl.devoteam.com\/working-at-devoteam\/\" target=\"_blank\" rel=\"noopener noreferrer\">Discover more<\/a><\/div>\n<\/div>\n","protected":false},"excerpt":{"rendered":"<p>\u00a0 NOTE: The Java example code for this technical blog can be found in this GitHub repo: GitHub Repo Consider the following situation: You want to use a single query to query a JDBC compatible database like Google Cloud SQL (MySQL) that contains millions of rows. You may want to do this when migrating legacy [&hellip;]<\/p>\n","protected":false},"featured_media":350997,"template":"","categories":[],"tags":[2297,1829,2717],"industry":[],"class_list":["post-601885","expert-view","type-expert-view","status-publish","has-post-thumbnail","hentry","tag-data-en-nl","tag-google-cloud-en-nl","tag-netherlands-en-nl"],"acf":[],"cards":"\n\t<div class=\"single-post-card\">\n\n\t\t<figure class=\"wp-block-post-featured-image\"><a href=\"https:\/\/devoteam.info\/en-nl\/expert-view\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\/\" target=\"_self\" ><img width=\"1000\" height=\"667\" src=\"https:\/\/devoteam.info\/wp-content\/uploads\/2024\/12\/Coding-image-general.jpeg\" class=\"attachment-post-thumbnail size-post-thumbnail wp-post-image\" alt=\"Querying JDBC database in parallel with Google Dataflow (Apache Beam)\" style=\"aspect-ratio:4\/3;width:100%;object-fit:cover;\" decoding=\"async\" loading=\"lazy\" srcset=\"https:\/\/devoteam.info\/wp-content\/uploads\/2024\/12\/Coding-image-general.jpeg 1000w, https:\/\/devoteam.info\/wp-content\/uploads\/2024\/12\/Coding-image-general-300x200.jpeg 300w, https:\/\/devoteam.info\/wp-content\/uploads\/2024\/12\/Coding-image-general-768x512.jpeg 768w\" sizes=\"auto, (max-width: 1000px) 100vw, 1000px\" \/><\/a><\/figure>\n\n\t\t\n\t\t<div class=\"wp-block-group is-vertical is-layout-flex wp-container-core-group-is-layout-43282307 wp-block-group-is-layout-flex\">\n\t<p style=\"font-style:normal;font-weight:700\" class=\"has-link-color wp-elements-1 wp-block-lp-post-type has-text-color has-primary-color has-small-font-size\">Expert View<\/p>\n\n\t\t\n\t\t<h3 style=\"font-style:normal;font-weight:400\" class=\"wp-block-post-title has-base-font-size\"><a href=\"https:\/\/devoteam.info\/en-nl\/expert-view\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\/\" target=\"_self\" >Querying JDBC database in parallel with Google Dataflow (Apache Beam)<\/a><\/h3><\/div>\n\t\t\n\t<\/div>\n\n","yoast_head":"<!-- This site is optimized with the Yoast SEO Premium plugin v28.4 (Yoast SEO v28.4) - https:\/\/yoast.com\/product\/yoast-seo-premium-wordpress\/ -->\n<title>Querying JDBC database in parallel with Google Dataflow (Apache Beam) | Devoteam<\/title>\n<meta name=\"robots\" content=\"index, follow, max-snippet:-1, max-image-preview:large, max-video-preview:-1\" \/>\n<link rel=\"canonical\" href=\"https:\/\/devoteam.info\/en-nl\/expert-view\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\/\" \/>\n<meta property=\"og:locale\" content=\"en_US\" \/>\n<meta property=\"og:type\" content=\"article\" \/>\n<meta property=\"og:title\" content=\"Querying JDBC database in parallel with Google Dataflow (Apache Beam)\" \/>\n<meta property=\"og:description\" content=\"\u00a0 NOTE: The Java example code for this technical blog can be found in this GitHub repo: GitHub Repo Consider the following situation: You want to use a single query to query a JDBC compatible database like Google Cloud SQL (MySQL) that contains millions of rows. You may want to do this when migrating legacy [&hellip;]\" \/>\n<meta property=\"og:url\" content=\"https:\/\/devoteam.info\/en-nl\/expert-view\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\/\" \/>\n<meta property=\"og:site_name\" content=\"Devoteam\" \/>\n<meta property=\"og:image\" content=\"https:\/\/devoteam.info\/wp-content\/uploads\/2024\/08\/Google-Cloud-logo.png\" \/>\n\t<meta property=\"og:image:width\" content=\"2048\" \/>\n\t<meta property=\"og:image:height\" content=\"362\" \/>\n\t<meta property=\"og:image:type\" content=\"image\/png\" \/>\n<meta name=\"twitter:card\" content=\"summary_large_image\" \/>\n<meta name=\"twitter:label1\" content=\"Est. reading time\" \/>\n\t<meta name=\"twitter:data1\" content=\"8 minutes\" \/>\n<script type=\"application\/ld+json\" class=\"yoast-schema-graph\">{\"@context\":\"https:\\\/\\\/schema.org\",\"@graph\":[{\"@type\":\"WebPage\",\"@id\":\"https:\\\/\\\/devoteam.info\\\/en-nl\\\/expert-view\\\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\\\/\",\"url\":\"https:\\\/\\\/devoteam.info\\\/en-nl\\\/expert-view\\\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\\\/\",\"name\":\"Querying JDBC database in parallel with Google Dataflow (Apache Beam) | Devoteam\",\"isPartOf\":{\"@id\":\"https:\\\/\\\/devoteam.info\\\/en-nl\\\/#website\"},\"primaryImageOfPage\":{\"@id\":\"https:\\\/\\\/devoteam.info\\\/en-nl\\\/expert-view\\\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\\\/#primaryimage\"},\"image\":{\"@id\":\"https:\\\/\\\/devoteam.info\\\/en-nl\\\/expert-view\\\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\\\/#primaryimage\"},\"thumbnailUrl\":\"https:\\\/\\\/devoteam.info\\\/wp-content\\\/uploads\\\/2024\\\/12\\\/Coding-image-general.jpeg\",\"datePublished\":\"2019-11-13T08:41:06+00:00\",\"breadcrumb\":{\"@id\":\"https:\\\/\\\/devoteam.info\\\/en-nl\\\/expert-view\\\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\\\/#breadcrumb\"},\"inLanguage\":\"en-NL\",\"potentialAction\":[{\"@type\":\"ReadAction\",\"target\":[\"https:\\\/\\\/devoteam.info\\\/en-nl\\\/expert-view\\\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\\\/\"]}]},{\"@type\":\"ImageObject\",\"inLanguage\":\"en-NL\",\"@id\":\"https:\\\/\\\/devoteam.info\\\/en-nl\\\/expert-view\\\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\\\/#primaryimage\",\"url\":\"https:\\\/\\\/devoteam.info\\\/wp-content\\\/uploads\\\/2024\\\/12\\\/Coding-image-general.jpeg\",\"contentUrl\":\"https:\\\/\\\/devoteam.info\\\/wp-content\\\/uploads\\\/2024\\\/12\\\/Coding-image-general.jpeg\",\"width\":1000,\"height\":667},{\"@type\":\"BreadcrumbList\",\"@id\":\"https:\\\/\\\/devoteam.info\\\/en-nl\\\/expert-view\\\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\\\/#breadcrumb\",\"itemListElement\":[{\"@type\":\"ListItem\",\"position\":1,\"name\":\"Home\",\"item\":\"https:\\\/\\\/devoteam.info\\\/en-nl\\\/\"},{\"@type\":\"ListItem\",\"position\":2,\"name\":\"Expert View\",\"item\":\"https:\\\/\\\/devoteam.info\\\/en-nl\\\/expert-view\\\/\"},{\"@type\":\"ListItem\",\"position\":3,\"name\":\"Querying JDBC database in parallel with Google Dataflow (Apache Beam)\"}]},{\"@type\":\"WebSite\",\"@id\":\"https:\\\/\\\/devoteam.info\\\/en-nl\\\/#website\",\"url\":\"https:\\\/\\\/devoteam.info\\\/en-nl\\\/\",\"name\":\"Devoteam\",\"description\":\"\",\"potentialAction\":[{\"@type\":\"SearchAction\",\"target\":{\"@type\":\"EntryPoint\",\"urlTemplate\":\"https:\\\/\\\/devoteam.info\\\/en-nl\\\/?s={search_term_string}\"},\"query-input\":{\"@type\":\"PropertyValueSpecification\",\"valueRequired\":true,\"valueName\":\"search_term_string\"}}],\"inLanguage\":\"en-NL\"}]}<\/script>\n<!-- \/ Yoast SEO Premium plugin. -->","yoast_head_json":{"title":"Querying JDBC database in parallel with Google Dataflow (Apache Beam) | Devoteam","robots":{"index":"index","follow":"follow","max-snippet":"max-snippet:-1","max-image-preview":"max-image-preview:large","max-video-preview":"max-video-preview:-1"},"canonical":"https:\/\/devoteam.info\/en-nl\/expert-view\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\/","og_locale":"en_US","og_type":"article","og_title":"Querying JDBC database in parallel with Google Dataflow (Apache Beam)","og_description":"\u00a0 NOTE: The Java example code for this technical blog can be found in this GitHub repo: GitHub Repo Consider the following situation: You want to use a single query to query a JDBC compatible database like Google Cloud SQL (MySQL) that contains millions of rows. You may want to do this when migrating legacy [&hellip;]","og_url":"https:\/\/devoteam.info\/en-nl\/expert-view\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\/","og_site_name":"Devoteam","og_image":[{"width":2048,"height":362,"url":"https:\/\/devoteam.info\/wp-content\/uploads\/2024\/08\/Google-Cloud-logo.png","type":"image\/png"}],"twitter_card":"summary_large_image","twitter_misc":{"Est. reading time":"8 minutes"},"schema":{"@context":"https:\/\/schema.org","@graph":[{"@type":"WebPage","@id":"https:\/\/devoteam.info\/en-nl\/expert-view\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\/","url":"https:\/\/devoteam.info\/en-nl\/expert-view\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\/","name":"Querying JDBC database in parallel with Google Dataflow (Apache Beam) | Devoteam","isPartOf":{"@id":"https:\/\/devoteam.info\/en-nl\/#website"},"primaryImageOfPage":{"@id":"https:\/\/devoteam.info\/en-nl\/expert-view\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\/#primaryimage"},"image":{"@id":"https:\/\/devoteam.info\/en-nl\/expert-view\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\/#primaryimage"},"thumbnailUrl":"https:\/\/devoteam.info\/wp-content\/uploads\/2024\/12\/Coding-image-general.jpeg","datePublished":"2019-11-13T08:41:06+00:00","breadcrumb":{"@id":"https:\/\/devoteam.info\/en-nl\/expert-view\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\/#breadcrumb"},"inLanguage":"en-NL","potentialAction":[{"@type":"ReadAction","target":["https:\/\/devoteam.info\/en-nl\/expert-view\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\/"]}]},{"@type":"ImageObject","inLanguage":"en-NL","@id":"https:\/\/devoteam.info\/en-nl\/expert-view\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\/#primaryimage","url":"https:\/\/devoteam.info\/wp-content\/uploads\/2024\/12\/Coding-image-general.jpeg","contentUrl":"https:\/\/devoteam.info\/wp-content\/uploads\/2024\/12\/Coding-image-general.jpeg","width":1000,"height":667},{"@type":"BreadcrumbList","@id":"https:\/\/devoteam.info\/en-nl\/expert-view\/querying-jdbc-database-in-parallel-with-google-dataflow-apache-beam\/#breadcrumb","itemListElement":[{"@type":"ListItem","position":1,"name":"Home","item":"https:\/\/devoteam.info\/en-nl\/"},{"@type":"ListItem","position":2,"name":"Expert View","item":"https:\/\/devoteam.info\/en-nl\/expert-view\/"},{"@type":"ListItem","position":3,"name":"Querying JDBC database in parallel with Google Dataflow (Apache Beam)"}]},{"@type":"WebSite","@id":"https:\/\/devoteam.info\/en-nl\/#website","url":"https:\/\/devoteam.info\/en-nl\/","name":"Devoteam","description":"","potentialAction":[{"@type":"SearchAction","target":{"@type":"EntryPoint","urlTemplate":"https:\/\/devoteam.info\/en-nl\/?s={search_term_string}"},"query-input":{"@type":"PropertyValueSpecification","valueRequired":true,"valueName":"search_term_string"}}],"inLanguage":"en-NL"}]}},"uagb_featured_image_src":{"full":["https:\/\/devoteam.info\/wp-content\/uploads\/2024\/12\/Coding-image-general.jpeg",1000,667,false],"thumbnail":["https:\/\/devoteam.info\/wp-content\/uploads\/2024\/12\/Coding-image-general-150x150.jpeg",150,150,true],"medium":["https:\/\/devoteam.info\/wp-content\/uploads\/2024\/12\/Coding-image-general-300x200.jpeg",300,200,true],"medium_large":["https:\/\/devoteam.info\/wp-content\/uploads\/2024\/12\/Coding-image-general-768x512.jpeg",768,512,true],"large":["https:\/\/devoteam.info\/wp-content\/uploads\/2024\/12\/Coding-image-general.jpeg",1000,667,false],"1536x1536":["https:\/\/devoteam.info\/wp-content\/uploads\/2024\/12\/Coding-image-general.jpeg",1000,667,false],"2048x2048":["https:\/\/devoteam.info\/wp-content\/uploads\/2024\/12\/Coding-image-general.jpeg",1000,667,false]},"uagb_author_info":{"display_name":"lea.mitteaux","author_link":"https:\/\/devoteam.info\/en-nl\/author\/"},"uagb_comment_info":0,"uagb_excerpt":"\u00a0 NOTE: The Java example code for this technical blog can be found in this GitHub repo: GitHub Repo Consider the following situation: You want to use a single query to query a JDBC compatible database like Google Cloud SQL (MySQL) that contains millions of rows. You may want to do this when migrating legacy&hellip;","_links":{"self":[{"href":"https:\/\/devoteam.info\/en-nl\/wp-json\/wp\/v2\/expert-view\/601885","targetHints":{"allow":["GET"]}}],"collection":[{"href":"https:\/\/devoteam.info\/en-nl\/wp-json\/wp\/v2\/expert-view"}],"about":[{"href":"https:\/\/devoteam.info\/en-nl\/wp-json\/wp\/v2\/types\/expert-view"}],"version-history":[{"count":0,"href":"https:\/\/devoteam.info\/en-nl\/wp-json\/wp\/v2\/expert-view\/601885\/revisions"}],"wp:featuredmedia":[{"embeddable":true,"href":"https:\/\/devoteam.info\/en-nl\/wp-json\/wp\/v2\/media\/350997"}],"wp:attachment":[{"href":"https:\/\/devoteam.info\/en-nl\/wp-json\/wp\/v2\/media?parent=601885"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"https:\/\/devoteam.info\/en-nl\/wp-json\/wp\/v2\/categories?post=601885"},{"taxonomy":"post_tag","embeddable":true,"href":"https:\/\/devoteam.info\/en-nl\/wp-json\/wp\/v2\/tags?post=601885"},{"taxonomy":"industry","embeddable":true,"href":"https:\/\/devoteam.info\/en-nl\/wp-json\/wp\/v2\/industry?post=601885"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}