[Bug]: JDBCIO Read without Partition occur GC overhead limit exceeded
- Dominant language
- Java
- Stars
- 8.7k
- Forks
- 4.7k
- Avg merge
- 1d 20h
- Merged PRs (30d)
- 196
Description
### What happened?
**My issue:**
When I read the mysql table using JDBCIO non-partitioned read mode, I get the following error:
```
Caused by: java.lang.OutOfMemoryError: GC overhead limit exceeded
at com.mysql.cj.protocol.a.NativePacketPayload.readBytes(NativePacketPayload.java:574)
at com.mysql.cj.protocol.a.NativePacketPayload.readBytes(NativePacketPayload.java:524)
at com.mysql.cj.protocol.a.TextRowFactory.createFromMessage(TextRowFactory.java:66)
at com.mysql.cj.protocol.a.TextRowFactory.createFromMessage(TextRowFactory.java:42)
at com.mysql.cj.protocol.a.ResultsetRowReader.read(ResultsetRowReader.java:87)
at com.mysql.cj.protocol.a.ResultsetRowReader.read(ResultsetRowReader.java:42)
at com.mysql.cj.protocol.a.NativeProtocol.read(NativeProtocol.java:1651)
at com.mysql.cj.protocol.a.TextResultsetReader.read(TextResultsetReader.java:87)
at com.mysql.cj.protocol.a.TextResultsetReader.read(TextResultsetReader.java:48)
at com.mysql.cj.protocol.a.NativeProtocol.read(NativeProtocol.java:1664)
at com.mysql.cj.protocol.a.NativeProtocol.readAllResults(NativeProtocol.java:1718)
at com.mysql.cj.protocol.a.NativeProtocol.sendQueryPacket(NativeProtocol.java:1064)
at com.mysql.cj.NativeSession.execSQL(NativeSession.java:665)
at com.mysql.cj.jdbc.ClientPreparedStatement.executeInternal(ClientPreparedStatement.java:893)
at com.mysql.cj.jdbc.ClientPreparedStatement.executeQuery(ClientPreparedStatement.java:972)
at org.apache.commons.dbcp2.DelegatingPreparedStatement.executeQuery(DelegatingPreparedStatement.java:122)
at org.apache.commons.dbcp2.DelegatingPreparedStatement.executeQuery(DelegatingPreparedStatement.java:122)
at org.apache.beam.sdk.io.jdbc.JdbcIO$ReadFn.processElement(JdbcIO.java:1399)
at org.apache.beam.sdk.io.jdbc.JdbcIO$ReadFn$DoFnInvoker.invokeProcessElement(Unknown Source)
at org.apache.beam.runners.core.SimpleDoFnRunner.invokeProcessElement(SimpleDoFnRunner.java:211)
at org.apache.beam.runners.core.SimpleDoFnRunner.processElement(SimpleDoFnRunner.java:188)
at org.apache.beam.runners.spark.translation.DoFnRunnerWithMetrics.processElement(DoFnRunnerWithMetrics.java:65)
at org.apache.beam.runners.spark.translation.SparkProcessContext$ProcCtxtIterator.computeNext(SparkProcessContext.java:140)
at org.apache.beam.vendor.guava.v26_0_jre.com.google.common.collect.AbstractIterator.tryToComputeNext(AbstractIterator.java:141)
at org.apache.beam.vendor.guava.v26_0_jre.com.google.common.collect.AbstractIterator.hasNext(AbstractIterator.java:136)
at scala.collection.convert.Wrappers$JIteratorWrapper.hasNext(Wrappers.scala:43)
at scala.collection.Iterator$$anon$12.hasNext(Iterator.scala:511)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:458)
at scala.collection.convert.Wrappers$IteratorWrapper.hasNext(Wrappers.scala:31)
at org.apache.beam.runners.spark.translation.MultiDoFnFunction.call(MultiDoFnFunction.java:128)
at org.apache.beam.runners.spark.translation.MultiDoFnFunction.call(MultiDoFnFunction.java:63)
at org.apache.spark.api.java.JavaRDDLike.$anonfun$mapPartitionsToPair$1(JavaRDDLike.scala:186)
```
Below code is trying to read mysql table:
```
public static PCollection ReadFromJdbc(Pipeline pipeline, JdbcConfig inputConfig, Schema inputSchema) {
LOG.info("read from jdbc, tablename is {}", inputConfig.getTableName());
String fieldNames = Joiner.on(",").join(inputSchema.getFieldNames());
return pipeline.apply(JdbcIO.read().withDataSourceConfiguration(
JdbcIO.DataSourceConfiguration.create(inputConfig.getDriverClassName(), inputConfig.getUrl())
.withUsername(inputConfig.getUsername()).withPassword(inputConfig.getPassword()))
.withQuery("select " + fieldNames + " from " + inputConfig.getTableName())
.withRowMapper(new JdbcIO.RowMapper() {
@Override
public Row mapRow(ResultSet resultSet) throws Exception {
Row.Builder builder = Row.withSchema(inputSchema);
for (int i = 1; i <= inputSchema.getFieldCount(); i++) {
Schema.TypeName type = inputSchema.getField(i - 1).getType().getTypeName();
switch (type) {
case DATETIME:
builder.addValue(new DateTime(resultSet.getTimestamp(i).getTime()));
continue;
default:
builder.addValue(resultSet.getObject(i));
}
}
return builder.build();
}
}).withCoder(SchemaCoder.of(inputSchema))
);
}
```
**Test table:**
10w records, 50 fields
**Spark enviroment:**
Alive Workers: 1
Cores in use: 15 Total, 0 Used
Memory in use: 29.0 GiB Total, 0.0 B Used
**Note: If I use JDBCIO read and set id to partition field, it works nice.**
### Issue Priority
Priority: 0 (outage / urgent vulnerability)
### Issue Components
- [ ] Component: Python SDK
- [x] Component: Java SDK
- [ ] Component: Go SDK
- [ ] Component: Typescript SDK
- [ ] Component: IO connector
- [ ] Component: Beam examples
- [ ] Component: Beam playground
- [ ] Component: Beam katas
- [ ] Component: Website
- [ ] Component: Spark Runner
- [ ] Component: Flink Runner
- [ ] Component: Samza Runner
- [ ] Component: Twister2 Runner
- [ ] Component: Hazelcast Jet Runner
- [ ] Component: Google Cloud Dataflow Runner
Contributor guide
Assessment
This issue has not been assessed yet.