diff --git a/CHANGES.txt b/CHANGES.txt index ac63fb3d3c..82f1d20755 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -43,6 +43,8 @@ 2.1.3 + * Invalidate affected prepared statements when a table's columns + are altered (CASSANDRA-7910) * Stress - user defined writes should populate sequentally (CASSANDRA-8524) * Fix regression in SSTableRewriter causing some rows to become unreadable during compaction (CASSANDRA-8429) diff --git a/src/java/org/apache/cassandra/config/CFMetaData.java b/src/java/org/apache/cassandra/config/CFMetaData.java index 0730ba745a..cb176f2880 100644 --- a/src/java/org/apache/cassandra/config/CFMetaData.java +++ b/src/java/org/apache/cassandra/config/CFMetaData.java @@ -724,11 +724,15 @@ public final class CFMetaData return def == null ? defaultValidator : def.type; } - public void reload() + /** + * Updates this object in place to match the definition in the system schema tables. + * @return true if any columns were added, removed, or altered; otherwise, false is returned + */ + public boolean reload() { try { - apply(LegacySchemaTables.createTableFromName(ksName, cfName)); + return apply(LegacySchemaTables.createTableFromName(ksName, cfName)); } catch (ConfigurationException e) { @@ -739,10 +743,11 @@ public final class CFMetaData /** * Updates CFMetaData in-place to match cfm * + * @return true if any columns were added, removed, or altered; otherwise, false is returned * @throws ConfigurationException if ks/cf names or cf ids didn't match */ @VisibleForTesting - public void apply(CFMetaData cfm) throws ConfigurationException + public boolean apply(CFMetaData cfm) throws ConfigurationException { logger.debug("applying {} to {}", cfm, this); @@ -802,6 +807,10 @@ public final class CFMetaData rebuild(); logger.debug("application result is {}", this); + + return !columnDiff.entriesOnlyOnLeft().isEmpty() || + !columnDiff.entriesOnlyOnRight().isEmpty() || + !columnDiff.entriesDiffering().isEmpty(); } public void validateCompatility(CFMetaData cfm) throws ConfigurationException diff --git a/src/java/org/apache/cassandra/config/Schema.java b/src/java/org/apache/cassandra/config/Schema.java index 21244ab4eb..694c05c23d 100644 --- a/src/java/org/apache/cassandra/config/Schema.java +++ b/src/java/org/apache/cassandra/config/Schema.java @@ -472,11 +472,11 @@ public class Schema { CFMetaData cfm = getCFMetaData(ksName, tableName); assert cfm != null; - cfm.reload(); + boolean columnsDidChange = cfm.reload(); Keyspace keyspace = Keyspace.open(cfm.ksName); keyspace.getColumnFamilyStore(cfm.cfName).reload(); - MigrationManager.instance.notifyUpdateColumnFamily(cfm); + MigrationManager.instance.notifyUpdateColumnFamily(cfm, columnsDidChange); } public void dropTable(String ksName, String tableName) diff --git a/src/java/org/apache/cassandra/cql3/QueryProcessor.java b/src/java/org/apache/cassandra/cql3/QueryProcessor.java index 8531d32574..f746a85c69 100644 --- a/src/java/org/apache/cassandra/cql3/QueryProcessor.java +++ b/src/java/org/apache/cassandra/cql3/QueryProcessor.java @@ -621,13 +621,22 @@ public class QueryProcessor implements QueryHandler } } + public void onUpdateColumnFamily(String ksName, String cfName, boolean columnsDidChange) + { + logger.info("Column definitions for {}.{} changed, invalidating related prepared statements", ksName, cfName); + if (columnsDidChange) + removeInvalidPreparedStatements(ksName, cfName); + } + public void onDropKeyspace(String ksName) { + logger.info("Keyspace {} was dropped, invalidating related prepared statements", ksName); removeInvalidPreparedStatements(ksName, null); } public void onDropColumnFamily(String ksName, String cfName) { + logger.info("Table {}.{} was dropped, invalidating related prepared statements", ksName, cfName); removeInvalidPreparedStatements(ksName, cfName); } diff --git a/src/java/org/apache/cassandra/service/MigrationListener.java b/src/java/org/apache/cassandra/service/MigrationListener.java index 2b728d9445..358b2367c7 100644 --- a/src/java/org/apache/cassandra/service/MigrationListener.java +++ b/src/java/org/apache/cassandra/service/MigrationListener.java @@ -47,7 +47,7 @@ public abstract class MigrationListener { } - public void onUpdateColumnFamily(String ksName, String cfName) + public void onUpdateColumnFamily(String ksName, String cfName, boolean columnsDidChange) { } diff --git a/src/java/org/apache/cassandra/service/MigrationManager.java b/src/java/org/apache/cassandra/service/MigrationManager.java index ef1adc6b2b..97f33f65c5 100644 --- a/src/java/org/apache/cassandra/service/MigrationManager.java +++ b/src/java/org/apache/cassandra/service/MigrationManager.java @@ -195,10 +195,10 @@ public class MigrationManager listener.onUpdateKeyspace(ksm.name); } - public void notifyUpdateColumnFamily(CFMetaData cfm) + public void notifyUpdateColumnFamily(CFMetaData cfm, boolean columnsDidChange) { for (MigrationListener listener : listeners) - listener.onUpdateColumnFamily(cfm.ksName, cfm.cfName); + listener.onUpdateColumnFamily(cfm.ksName, cfm.cfName, columnsDidChange); } public void notifyUpdateUserType(UserType ut) diff --git a/src/java/org/apache/cassandra/transport/Server.java b/src/java/org/apache/cassandra/transport/Server.java index 147d7292ba..d1fc74470e 100644 --- a/src/java/org/apache/cassandra/transport/Server.java +++ b/src/java/org/apache/cassandra/transport/Server.java @@ -429,7 +429,7 @@ public class Server implements CassandraDaemon.Server server.connectionTracker.send(new Event.SchemaChange(Event.SchemaChange.Change.UPDATED, ksName)); } - public void onUpdateColumnFamily(String ksName, String cfName) + public void onUpdateColumnFamily(String ksName, String cfName, boolean columnsDidChange) { server.connectionTracker.send(new Event.SchemaChange(Event.SchemaChange.Change.UPDATED, Event.SchemaChange.Target.TABLE, ksName, cfName)); }