diff --git a/snd/src/org/labkey/snd/SNDManager.java b/snd/src/org/labkey/snd/SNDManager.java index b6f37af5..7826ce29 100644 --- a/snd/src/org/labkey/snd/SNDManager.java +++ b/snd/src/org/labkey/snd/SNDManager.java @@ -91,6 +91,7 @@ import org.labkey.snd.security.SNDSecurityManager; import org.labkey.snd.trigger.SNDTriggerManager; +import java.nio.ByteBuffer; import java.sql.SQLException; import java.text.ParseException; import java.text.SimpleDateFormat; @@ -104,7 +105,9 @@ import java.util.HashSet; import java.util.Iterator; import java.util.List; +import java.util.LongSummaryStatistics; import java.util.Map; +import java.util.Objects; import java.util.Optional; import java.util.Set; import java.util.TreeMap; @@ -153,6 +156,14 @@ public static UserSchema getSndUserSchemaAdminRole(Container c, User u) public static int MAX_MERGE_ROWS = 2000; + /** Below the ETL batch size so that the full batches of an initial load suppress the id lists while the smaller batches of an incremental run keep them. */ + public static final int MAX_LOGGED_IDS = 2000; + private static final int LOGGED_IDS_PER_LINE = 250; + private static final int LOGGED_PAIRS_PER_LINE = 50; + + /** The incremental filter column of both SND ETL source views. */ + public static final String SOURCE_ROWVERSION_COLUMN = "timestamp"; + public static Logger getLogger(Map configParameters, Class clazz) { Logger log = null; @@ -164,6 +175,87 @@ public static Logger getLogger(Map configParameters, Class claz return log; } + /** + * Writes an id set to the ETL job log so that the id sets logged by different ETL steps of the same run can be + * diffed against each other. Chunked because a single line of thousands of ids is unreadable, and capped because + * the initial full data load would otherwise write the entire table to the log. + */ + public static void logIds(Logger log, String message, Collection ids) + { + if (!log.isDebugEnabled()) + return; + + log.debug(message + " Count: " + ids.size() + "."); + + if (ids.isEmpty()) + return; + + if (ids.size() > MAX_LOGGED_IDS) + { + log.debug("Id list omitted, more than " + MAX_LOGGED_IDS + " ids."); + return; + } + + List sorted = ids.stream().filter(Objects::nonNull).sorted().collect(Collectors.toList()); + for (List chunk : ListUtils.partition(sorted, LOGGED_IDS_PER_LINE)) + log.debug(" " + StringUtils.join(chunk, ", ")); + } + + /** + * SQL Server hands a rowversion back as binary(8); the ETL's own persisted window state returns it as a number. + * Read big-endian, matching how the incremental filter logs its bounds, so the two can be compared directly. + */ + @Nullable + public static Long toRowversion(@Nullable Object o) + { + if (o instanceof byte[] bytes && 8 == bytes.length) + return ByteBuffer.wrap(bytes).getLong(); + if (o instanceof Number n) + return n.longValue(); + return null; + } + + /** + * Logs the rowversion span of a batch so it can be placed against the incremental window the ETL logged for the + * run. Both SND source views draw their rowversions from the same source database, so the spans the two steps + * report are on one sequence and comparable. + */ + public static void logRowversionRange(Logger log, String message, Collection> rows) + { + if (!log.isDebugEnabled()) + return; + + LongSummaryStatistics stats = rows.stream() + .map(row -> toRowversion(row.get(SOURCE_ROWVERSION_COLUMN))) + .filter(Objects::nonNull) + .mapToLong(Long::longValue) + .summaryStatistics(); + + if (0 == stats.getCount()) + log.debug(message + " No source rowversions in this batch."); + else + log.debug(message + " Rowversions " + stats.getMin() + " to " + stats.getMax() + " over " + stats.getCount() + " rows."); + } + + /** + * Pairs each id with its source rowversion. Logged separately from the bare list the same set gets from logIds, + * which stays free of annotations so it can be diffed against the other step's list. + */ + public static void logIdRowversions(Logger log, String message, Collection ids, Map rowversions) + { + if (!log.isDebugEnabled() || ids.isEmpty() || ids.size() > MAX_LOGGED_IDS) + return; + + log.debug(message); + + List pairs = ids.stream().filter(Objects::nonNull).sorted() + .map(id -> id + ":" + rowversions.get(id)) + .collect(Collectors.toList()); + + for (List chunk : ListUtils.partition(pairs, LOGGED_PAIRS_PER_LINE)) + log.debug(" " + StringUtils.join(chunk, ", ")); + } + public static String getPackageName(int id) { return PackageDomainKind.getPackageKindName() + "-" + id; diff --git a/snd/src/org/labkey/snd/query/AttributeDataTable.java b/snd/src/org/labkey/snd/query/AttributeDataTable.java index 81d34367..cb4a490d 100644 --- a/snd/src/org/labkey/snd/query/AttributeDataTable.java +++ b/snd/src/org/labkey/snd/query/AttributeDataTable.java @@ -149,6 +149,10 @@ public QueryUpdateService getUpdateService() protected class UpdateService extends SNDQueryUpdateService { + /** Bounds the source ordering check below. It costs one retained URI per distinct EventDataId, and an ungrouped source would otherwise log once per row. */ + private static final int MAX_TRACKED_URIS = 50_000; + private static final int MAX_ORDER_WARNINGS = 10; + private final SNDManager _sndManager = SNDManager.get(); private final SNDService _sndService = SNDService.get(); private final DbSchema _expSchema = OntologyManager.getExpSchema(); @@ -232,6 +236,26 @@ private List> updateObjectProperty(User user, Container cont { logger.info("Begin updating exp.ObjectProperty."); + // An EventDataId gets one source row per attribute; keep the newest, since that is the one that pulled it into the window. + boolean trackRowversions = logger.isDebugEnabled(); + Set incomingEventDataIds = new HashSet<>(); + Map rowversionByEventDataId = new HashMap<>(); + for (Map row : data) + { + Integer eventDataId = (Integer) row.get("EventDataId"); + incomingEventDataIds.add(eventDataId); + + if (trackRowversions) + { + Long rowversion = SNDManager.toRowversion(row.get(SNDManager.SOURCE_ROWVERSION_COLUMN)); + if (null != rowversion) + rowversionByEventDataId.merge(eventDataId, rowversion, Math::max); + } + } + + SNDManager.logIds(logger, "Source rows: " + data.size() + ". EventDataIds in this batch:", incomingEventDataIds); + SNDManager.logRowversionRange(logger, "Source span of this batch.", data); + int inserted = 0; String prevUri = null; @@ -240,6 +264,10 @@ private List> updateObjectProperty(User user, Container cont boolean found = false; Set cacheEventIds = new HashSet<>(); + Set writtenEventDataIds = new HashSet<>(); + Set flushedUris = new HashSet<>(); + boolean checkOrdering = logger.isDebugEnabled(); + int outOfOrderFlushes = 0; for(Map row : data) { @@ -255,7 +283,8 @@ private List> updateObjectProperty(User user, Container cont //add to list of cached narrative rows to delete cacheEventIds.add((Integer) row.get("EventId")); - String objectURI = getObjectURI((Integer) row.get("EventDataId"), container); + Integer eventDataId = (Integer) row.get("EventDataId"); + String objectURI = getObjectURI(eventDataId, container); if (prevUri == null) prevUri = objectURI; @@ -292,11 +321,11 @@ else if (stringValue != null) { if (pd.getLookupSchema() != null && pd.getLookupQuery() != null) { - logger.info("Value null for property " + pd.getName() + ". Value skipped. Verify lookup " + pd.getLookupSchema() + "." + pd.getLookupQuery() + " contains " + stringValue); + logger.info("Value null for property " + pd.getName() + ", EventDataId: " + eventDataId + ". Value skipped. Verify lookup " + pd.getLookupSchema() + "." + pd.getLookupQuery() + " contains " + stringValue); } else { - logger.info("Value null for property " + pd.getName() + ". Value skipped."); + logger.info("Value null for property " + pd.getName() + ", EventDataId: " + eventDataId + ". Value skipped."); } } @@ -320,10 +349,27 @@ else if (stringValue != null) if (!prevUri.equals(objectURI)) { inserted = insertObject(container, user, prevUri, prevObjProps, pkgId, inserted, logger); + + // Properties are only flushed when the URI changes, so a URI seen twice means the source + // did not arrive grouped by EventDataId and the ORDER BY in v_snd_attributeData was lost. + if (checkOrdering) + { + if (!flushedUris.add(prevUri) && ++outOfOrderFlushes <= MAX_ORDER_WARNINGS) + logger.debug("Source rows are not grouped by EventDataId; exp.ObjectProperty for {} was written in more than one pass.", prevUri); + + if (flushedUris.size() >= MAX_TRACKED_URIS) + { + logger.debug("More than {} EventDataIds in this batch; ending the source ordering check.", MAX_TRACKED_URIS); + flushedUris.clear(); + checkOrdering = false; + } + } + prevUri = objectURI; prevObjProps = new ArrayList<>(); } prevObjProps.add(oprop); + writtenEventDataIds.add(eventDataId); } } @@ -332,12 +378,16 @@ else if (stringValue != null) } if (!found) { - throw new RuntimeException("Attribute metadata not found for key: '" + key + "' in package: " + pkgId); + throw new RuntimeException("Attribute metadata not found for key: '" + key + "' in package: " + pkgId + + ", EventDataId: " + eventDataId + ". Aborting, leaving all " + incomingEventDataIds.size() + + " EventDataIds in this batch with the attribute values the _SND Event Data step already cleared."); } } else { - throw new RuntimeException("Package metadata not found for package id: " + pkgId); + throw new RuntimeException("Package metadata not found for package id: " + pkgId + + ", EventDataId: " + eventDataId + ". Aborting, leaving all " + incomingEventDataIds.size() + + " EventDataIds in this batch with the attribute values the _SND Event Data step already cleared."); } } @@ -347,8 +397,28 @@ else if (stringValue != null) } OntologyManager.clearPropertyCache(); + logger.info("End updating exp.ObjectProperty. Inserted/Updated " + inserted + " rows."); + SNDManager.logIds(logger, "EventDataIds written:", writtenEventDataIds); + + if (outOfOrderFlushes > MAX_ORDER_WARNINGS) + logger.debug("{} objectURIs in total were written in more than one pass; further messages were suppressed.", outOfOrderFlushes); + + // Collect only the misses; copying the incoming set would double its footprint on a full load. + Set unwritten = new HashSet<>(); + for (Integer id : incomingEventDataIds) + { + if (!writtenEventDataIds.contains(id)) + unwritten.add(id); + } + + if (!unwritten.isEmpty()) + { + SNDManager.logIds(logger, "EventDataIds present in the source rows but left with no attribute values written:", unwritten); + SNDManager.logIdRowversions(logger, "Rowversions of those EventDataIds, to place them against the incremental window of this run:", unwritten, rowversionByEventDataId); + } + _sndManager.updateNarrativeCache(container, user, cacheEventIds, logger); return data; diff --git a/snd/src/org/labkey/snd/query/EventDataTable.java b/snd/src/org/labkey/snd/query/EventDataTable.java index 1c9bce9d..57491ebd 100644 --- a/snd/src/org/labkey/snd/query/EventDataTable.java +++ b/snd/src/org/labkey/snd/query/EventDataTable.java @@ -25,6 +25,7 @@ import org.labkey.api.data.JdbcType; import org.labkey.api.data.SQLFragment; import org.labkey.api.data.SqlExecutor; +import org.labkey.api.data.SqlSelector; import org.labkey.api.data.TableInfo; import org.labkey.api.dataiterator.DataIteratorBuilder; import org.labkey.api.dataiterator.DataIteratorContext; @@ -48,7 +49,9 @@ import java.io.IOException; import java.sql.SQLException; +import java.util.ArrayList; import java.util.HashSet; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Set; @@ -108,6 +111,9 @@ public QueryUpdateService getUpdateService() protected static class UpdateService extends SNDQueryUpdateService { + /** Keeps the ObjectURI IN clause well under the SQL Server parameter limit. */ + private static final int URI_CHUNK_SIZE = 500; + private final SNDManager _sndManager = SNDManager.get(); private final SNDService _sndService = SNDService.get(); private final DbSchema _expSchema = OntologyManager.getExpSchema(); @@ -122,6 +128,62 @@ private String getObjectURI(Integer eventDataId, Container c) return _sndManager.generateLsid(c, String.valueOf(eventDataId)); } + /** + * EventDataIds in this batch whose exp.Object currently carries attribute values. Deleting the exp.Object row + * cascades to exp.ObjectProperty, so these are the values the merge destroys; only the _SND Attribute Data ETL + * step re-inserts them, and it computes its incremental window independently of this step's. + */ + private Set getEventDataIdsWithAttributeData(Container container, Map eventDataIdsByUri) + { + Set withAttributeData = new HashSet<>(); + List uris = new ArrayList<>(eventDataIdsByUri.keySet()); + + for (int i = 0; i < uris.size(); i += URI_CHUNK_SIZE) + { + List chunk = uris.subList(i, Math.min(i + URI_CHUNK_SIZE, uris.size())); + + // EXISTS rather than a join: both indexes (UQ_Object on ObjectURI, PK_ObjectProperty on ObjectId) + // are seeks, and the semi-join stops at the first property instead of reading all of them per object. + SQLFragment sql = new SQLFragment("SELECT o.ObjectURI FROM ") + .append(OntologyManager.getTinfoObject(), "o") + .append(" WHERE o.Container = ?").add(container.getId()) + .append(" AND EXISTS (SELECT 1 FROM ").append(OntologyManager.getTinfoObjectProperty(), "op") + .append(" WHERE op.ObjectId = o.ObjectId)") + .append(" AND o.ObjectURI").appendInClause(chunk, _expSchema.getSqlDialect()); + + new SqlSelector(_expSchema, sql).getCollection(String.class) + .forEach(uri -> withAttributeData.add(eventDataIdsByUri.get(uri))); + } + + return withAttributeData; + } + + /** + * Diagnostic only, so a failure here must not abort the merge. Skipped above the cap logIds lists at, where the + * chunked queries would cost hundreds of round trips to produce a bare count. + */ + private void logAttributeDataToBeCleared(Container container, Map eventDataIdsByUri, Logger log) + { + if (!log.isDebugEnabled()) + return; + + if (eventDataIdsByUri.size() > SNDManager.MAX_LOGGED_IDS) + { + log.debug("More than " + SNDManager.MAX_LOGGED_IDS + " EventDataIds in this batch; skipping the check for attribute values about to be cleared."); + return; + } + + try + { + SNDManager.logIds(log, "Attribute values about to be cleared by this merge; the _SND Attribute Data step must re-insert them.", + getEventDataIdsWithAttributeData(container, eventDataIdsByUri)); + } + catch (Exception e) + { + log.debug("Could not determine which EventDataIds have attribute values; continuing with the merge.", e); + } + } + @Override public int mergeRows(User user, Container container, DataIteratorBuilder rows, BatchValidationException errors, @Nullable Map configParameters, Map extraScriptContext) @@ -158,14 +220,28 @@ public int mergeRows(User user, Container container, DataIteratorBuilder rows, B log.info("Merging rows."); log.info("Begin updating exp.Object table."); - int count = 0; - for(Map map : data) + + Map eventDataIdsByUri = new LinkedHashMap<>(); + for (Map map : data) { - String objectURI = getObjectURI((Integer) map.get("EventDataId"), container); + Integer eventDataId = (Integer) map.get("EventDataId"); + String objectURI = getObjectURI(eventDataId, container); //update snd.EventData row with objectURI map.put("ObjectURI", objectURI); + eventDataIdsByUri.put(objectURI, eventDataId); + } + + SNDManager.logIds(log, "EventDataIds merged into snd.EventData by this batch:", eventDataIdsByUri.values()); + SNDManager.logRowversionRange(log, "Source span of the merged rows.", data); + logAttributeDataToBeCleared(container, eventDataIdsByUri, log); + + int count = 0; + for(Map map : data) + { + String objectURI = (String) map.get("ObjectURI"); + //delete row from exp.Object OntologyManager.deleteOntologyObjects(container, objectURI); @@ -215,9 +291,11 @@ public int importRows(User user, Container container, DataIteratorBuilder rows, log.info("Begin inserting into exp.Object."); int count = 0; + Set eventDataIds = new HashSet<>(); for(Map map : data) { - String objectURI = getObjectURI((Integer) map.get("EventDataId"), container); + Integer eventDataId = (Integer) map.get("EventDataId"); + String objectURI = getObjectURI(eventDataId, container); //update snd.EventData row with objectURI map.put("ObjectURI", objectURI); @@ -228,6 +306,8 @@ public int importRows(User user, Container container, DataIteratorBuilder rows, //add to list of cached narrative rows to delete cacheData.add((Integer) map.get("EventId")); + eventDataIds.add(eventDataId); + count++; //TODO: Count in exp.Object is not going to be the same as in snd.EventData - need to figure out how to get the count to log if(count % 1000 == 0) @@ -235,6 +315,10 @@ public int importRows(User user, Container container, DataIteratorBuilder rows, } log.info("End inserting into exp.Object. Inserted total of " + count + " rows."); + // These rows get a fresh exp.Object with no properties, so they depend on the _SND Attribute Data step just as much as the merged ones do. + SNDManager.logIds(log, "EventDataIds inserted into snd.EventData by this batch:", eventDataIds); + SNDManager.logRowversionRange(log, "Source span of the inserted rows.", data); + DataIteratorBuilder rowsWithObjectURI = new ListofMapsDataIterator.Builder(data.get(0).keySet(), data); _sndManager.updateNarrativeCache(container, user, cacheData, log); @@ -334,11 +418,15 @@ private void deleteFromExpTables(List> oldRows, Container co { log.info("Begin deleting from exp.ObjectProperty and exp.Object."); int count = 0; + Set eventDataIds = new HashSet<>(); //This will be a cascading delete across exp.ObjectProperty, exp.Object, and snd.EventData for (Map map : oldRows) { - String objectURI = getObjectURI((Integer) map.get("EventDataId"), container); + Integer eventDataId = (Integer) map.get("EventDataId"); + String objectURI = getObjectURI(eventDataId, container); + + eventDataIds.add(eventDataId); OntologyObject obj = OntologyManager.getOntologyObject(container, objectURI); //delete row from exp.ObjectProperty @@ -355,6 +443,9 @@ private void deleteFromExpTables(List> oldRows, Container co } log.info("End deleting from exp.ObjectProperty and exp.Object. Deleted total of " + count + " rows."); + + // Without these the deleted rows read as attribute data the _SND Attribute Data step failed to write. + SNDManager.logIds(log, "EventDataIds deleted from snd.EventData by this batch:", eventDataIds); } private int deleteAllFromExpTables(Logger log)