From f56e3815736dbbe5c156a82cf2ea25f09a9946c1 Mon Sep 17 00:00:00 2001 From: Viraj Jasani Date: Thu, 29 Oct 2020 15:27:28 +0530 Subject: [PATCH] PHOENIX-6002 : Resolve connection leak through QueryUtil.getConnectionOnServer() --- .../IndexHalfStoreFileReaderGenerator.java | 16 ++++----- .../coprocessor/DropColumnMutator.java | 34 +++++++++++-------- .../coprocessor/MetaDataEndpointImpl.java | 24 ++++++++----- .../mapreduce/FormatToKeyValueReducer.java | 4 +-- 4 files changed, 43 insertions(+), 35 deletions(-) diff --git a/phoenix-core/src/main/java/org/apache/hadoop/hbase/regionserver/IndexHalfStoreFileReaderGenerator.java b/phoenix-core/src/main/java/org/apache/hadoop/hbase/regionserver/IndexHalfStoreFileReaderGenerator.java index 7254f573fa6..5e591361680 100644 --- a/phoenix-core/src/main/java/org/apache/hadoop/hbase/regionserver/IndexHalfStoreFileReaderGenerator.java +++ b/phoenix-core/src/main/java/org/apache/hadoop/hbase/regionserver/IndexHalfStoreFileReaderGenerator.java @@ -276,24 +276,22 @@ private InternalScanner getRepairScanner(RegionCoprocessorEnvironment env, Store scan.addFamily(s.getColumnFamilyDescriptor().getName()); } } - try { - PhoenixConnection conn = QueryUtil.getConnectionOnServer(env.getConfiguration()) - .unwrap(PhoenixConnection.class); + try (PhoenixConnection conn = QueryUtil.getConnectionOnServer( + env.getConfiguration()).unwrap(PhoenixConnection.class)) { PTable dataPTable = IndexUtil.getPDataTable(conn, env.getRegion().getTableDescriptor()); final List maintainers = Lists - .newArrayListWithExpectedSize(dataPTable.getIndexes().size()); + .newArrayListWithExpectedSize(dataPTable.getIndexes().size()); for (PTable index : dataPTable.getIndexes()) { if (index.getIndexType() == IndexType.LOCAL) { maintainers.add(index.getIndexMaintainer(dataPTable, conn)); } } - return new DataTableLocalIndexRegionScanner(env.getRegion().getScanner(scan), env.getRegion(), - maintainers, store.getColumnFamilyDescriptor().getName(),env.getConfiguration()); - - + return new DataTableLocalIndexRegionScanner( + env.getRegion().getScanner(scan), env.getRegion(), maintainers, + store.getColumnFamilyDescriptor().getName(), + env.getConfiguration()); } catch (SQLException e) { throw new IOException(e); - } } } diff --git a/phoenix-core/src/main/java/org/apache/phoenix/coprocessor/DropColumnMutator.java b/phoenix-core/src/main/java/org/apache/phoenix/coprocessor/DropColumnMutator.java index f1491b94aca..1fe9d847a0b 100644 --- a/phoenix-core/src/main/java/org/apache/phoenix/coprocessor/DropColumnMutator.java +++ b/phoenix-core/src/main/java/org/apache/phoenix/coprocessor/DropColumnMutator.java @@ -132,23 +132,27 @@ public MetaDataMutationResult validateWithChildViews(PTable table, List if (existingViewColumn != null && view.getViewStatement() != null) { ParseNode viewWhere = new SQLParser(view.getViewStatement()).parseQuery().getWhere(); - PhoenixConnection conn = QueryUtil.getConnectionOnServer(conf).unwrap( - PhoenixConnection.class); - PhoenixStatement statement = new PhoenixStatement(conn); - TableRef baseTableRef = new TableRef(view); - ColumnResolver columnResolver = FromCompiler.getResolver(baseTableRef); - StatementContext context = new StatementContext(statement, columnResolver); - Expression whereExpression = WhereCompiler.compile(context, viewWhere); - Expression colExpression = - new ColumnRef(baseTableRef, existingViewColumn.getPosition()) - .newColumnExpression(); - MetaDataEndpointImpl.ColumnFinder columnFinder = - new MetaDataEndpointImpl.ColumnFinder(colExpression); - whereExpression.accept(columnFinder); - if (columnFinder.getColumnFound()) { - return new MetaDataProtocol.MetaDataMutationResult( + try (PhoenixConnection conn = + QueryUtil.getConnectionOnServer(conf) + .unwrap(PhoenixConnection.class)) { + PhoenixStatement statement = new PhoenixStatement(conn); + TableRef baseTableRef = new TableRef(view); + ColumnResolver columnResolver = + FromCompiler.getResolver(baseTableRef); + StatementContext context = + new StatementContext(statement, columnResolver); + Expression whereExpression = + WhereCompiler.compile(context, viewWhere); + Expression colExpression = new ColumnRef(baseTableRef, + existingViewColumn.getPosition()).newColumnExpression(); + MetaDataEndpointImpl.ColumnFinder columnFinder = + new MetaDataEndpointImpl.ColumnFinder(colExpression); + whereExpression.accept(columnFinder); + if (columnFinder.getColumnFound()) { + return new MetaDataProtocol.MetaDataMutationResult( MetaDataProtocol.MutationCode.UNALLOWED_TABLE_MUTATION, EnvironmentEdgeManager.currentTimeMillis(), table); + } } } diff --git a/phoenix-core/src/main/java/org/apache/phoenix/coprocessor/MetaDataEndpointImpl.java b/phoenix-core/src/main/java/org/apache/phoenix/coprocessor/MetaDataEndpointImpl.java index a830d7fd55e..e4e674210cc 100644 --- a/phoenix-core/src/main/java/org/apache/phoenix/coprocessor/MetaDataEndpointImpl.java +++ b/phoenix-core/src/main/java/org/apache/phoenix/coprocessor/MetaDataEndpointImpl.java @@ -2290,16 +2290,16 @@ public void dropTable(RpcController controller, DropTableRequest request, && (clientVersion >= MIN_SPLITTABLE_SYSTEM_CATALOG || SchemaUtil.getPhysicalTableName(SYSTEM_CHILD_LINK_NAME_BYTES, env.getConfiguration()).equals(hTable.getName()))) { - try { - PhoenixConnection conn = - QueryUtil.getConnectionOnServer(env.getConfiguration()) - .unwrap(PhoenixConnection.class); + try (PhoenixConnection conn = + QueryUtil.getConnectionOnServer(env.getConfiguration()) + .unwrap(PhoenixConnection.class)) { Task.addTask(conn, PTable.TaskType.DROP_CHILD_VIEWS, - Bytes.toString(tenantIdBytes), Bytes.toString(schemaName), - Bytes.toString(tableOrViewName), - PTable.TaskStatus.CREATED.toString(), - null, null, null, null, - this.accessCheckEnabled); + Bytes.toString(tenantIdBytes), + Bytes.toString(schemaName), + Bytes.toString(tableOrViewName), + PTable.TaskStatus.CREATED.toString(), + null, null, null, null, + this.accessCheckEnabled); } catch (Throwable t) { LOGGER.error("Adding a task to drop child views failed!", t); } @@ -3201,6 +3201,9 @@ else if (isCoveredColumn) { invalidateList.add(new ImmutableBytesPtr(indexKey)); } } + if (connection != null) { + connection.close(); + } return null; } @@ -3268,6 +3271,9 @@ else if (isCoveredColumn) { index.getSchemaName().getBytes() : ByteUtil.EMPTY_BYTE_ARRAY, index.getTableName().getBytes()); } } + if (connection != null) { + connection.close(); + } return null; } diff --git a/phoenix-core/src/main/java/org/apache/phoenix/mapreduce/FormatToKeyValueReducer.java b/phoenix-core/src/main/java/org/apache/phoenix/mapreduce/FormatToKeyValueReducer.java index ea727f7e998..9dccd6c3403 100644 --- a/phoenix-core/src/main/java/org/apache/phoenix/mapreduce/FormatToKeyValueReducer.java +++ b/phoenix-core/src/main/java/org/apache/phoenix/mapreduce/FormatToKeyValueReducer.java @@ -77,8 +77,8 @@ protected void setup(Context context) throws IOException, InterruptedException { for (Map.Entry entry : conf) { clientInfos.setProperty(entry.getKey(), entry.getValue()); } - try { - PhoenixConnection conn = (PhoenixConnection) QueryUtil.getConnectionOnServer(clientInfos, conf); + try (PhoenixConnection conn = (PhoenixConnection) QueryUtil + .getConnectionOnServer(clientInfos, conf)) { builder = conn.getKeyValueBuilder(); final String tableNamesConf = conf.get(FormatToBytesWritableMapper.TABLE_NAMES_CONFKEY); final String logicalNamesConf = conf.get(FormatToBytesWritableMapper.LOGICAL_NAMES_CONFKEY);