Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions importers/build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,8 @@ dependencies {
shadowJar {
archiveBaseName = "xlimporter"
archiveClassifier = null // removes `-all` in the filename of the created .jar
duplicatesStrategy = DuplicatesStrategy.INCLUDE
mergeServiceFiles()
}

jar {
Expand Down
72 changes: 50 additions & 22 deletions importers/src/main/groovy/whelk/importer/DatasetImporter.groovy
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import whelk.JsonLd
import whelk.TargetVocabMapper
import whelk.Whelk
import whelk.converter.TrigToJsonLdParser
import whelk.converter.JsonLdShapes
import whelk.util.DocumentUtil

import java.time.Duration
Expand Down Expand Up @@ -77,9 +78,20 @@ class DatasetImporter {
DatasetImporter(Whelk whelk, String datasetUri, Map flags=[:], Object descriptions=null) {
this.whelk = whelk
this.datasetUri = datasetUri
if (datasetUri != null) {
log.info("Initialized DatasetImporter for ${datasetUri}")
}
log.info("Using system context: ${whelk.systemContextUri}")
if (whelk.systemContextUri) {
contextDocData = getDocByMainEntityId(whelk.systemContextUri)?.data
if (contextDocData.containsKey(CONTEXT)) {
def ctx = contextDocData.get(CONTEXT)
if (ctx instanceof Map) log.info("Context size: ${ctx.size()}")
} else {
log.warn("Context missing ${CONTEXT}")
}
}

if (descriptions != null) {
Map datasetDesc = descriptions instanceof Map ? (Map) descriptions : loadData((String) descriptions)
givenDsData = (Map) findInData(datasetDesc, datasetUri)
Expand All @@ -96,13 +108,14 @@ class DatasetImporter {
}

static void loadDescribedDatasets(Whelk whelk, String datasetDescPath, String sourceBaseDir, Set<String> onlyDatasets=null, Map flags=[:]) {
log.info("Loading datasets described in: ${datasetDescPath}")
var dsImp = new DatasetImporter(whelk, null)
var datasets = (Map) new File(datasetDescPath).withInputStream {
dsImp.contextDocData ? dsImp.loadTurtleAsSystemShaped(it) : loadSelfCompactedTurtle(it)
}
for (Map item : (List<Map>) datasets[GRAPH] ?: asList(datasets)) {
if (onlyDatasets && item[ID] !in onlyDatasets) {
System.err.println("Skipping dataset: ${item[ID]}")
log.info("Skipping dataset: ${item[ID]}")
continue
}
if (item[TYPE] == 'Dataset' && 'sourceData' in item) {
Expand Down Expand Up @@ -140,7 +153,7 @@ class DatasetImporter {

private void doImportDataset(String sourceUrl) {
long startTime = System.nanoTime()
System.err.println("Importing from: ${sourceUrl}")
log.info("Importing from: ${sourceUrl}")

Set<String> idsInInput = []

Expand Down Expand Up @@ -185,7 +198,7 @@ class DatasetImporter {
countWriteAndMaybeFlushIndexing()

if ( lineCount % 100 == 0 ) {
System.err.println("Processed " + lineCount + " input records. " + createdCount + " created, " +
log.info("Processed " + lineCount + " input records. " + createdCount + " created, " +
updatedCount + " updated, " + (lineCount-createdCount-updatedCount) + " already up to date.")
}
++lineCount
Expand All @@ -205,11 +218,11 @@ class DatasetImporter {

Duration elapsedTime = Duration.ofNanos(System.nanoTime() - startTime)
String elapsed = String.format("%02dh%02dm%02ds", elapsedTime.toHours(), elapsedTime.toMinutesPart(), elapsedTime.toSecondsPart())
System.err.println("Created: " + createdCount +" new,\n" +
"updated: " + updatedCount + " existing and\n" +
"deleted: " + deletedCount + " old records (should have been: " + (deletedCount + needsRetry.size()) + "),\n" +
"out of the: " + idsInInput.size() + " records in dataset: \"" + dsInfo.uri + "\".\n" +
"Dataset now in sync in ${elapsed}.")
log.info("Created: " + createdCount +" new,\n" +
"\tupdated: " + updatedCount + " existing and\n" +
"\tdeleted: " + deletedCount + " old records (should have been: " + (deletedCount + needsRetry.size()) + "),\n" +
"\tout of the: " + idsInInput.size() + " records in dataset: \"" + dsInfo.uri + "\".\n" +
"\tDataset now in sync in ${elapsed}.")
}

void dropDataset() {
Expand All @@ -223,7 +236,7 @@ class DatasetImporter {
} finally {
whelk.endDeferredIndexing()
}
System.err.println("Deleted dataset ${dsInfo.uri} with ${deletedCount} existing records")
log.info("Deleted dataset ${dsInfo.uri} with ${deletedCount} existing records")
}

private void processDataset(String sourceUrl, Closure processItem) {
Expand Down Expand Up @@ -251,16 +264,16 @@ class DatasetImporter {
Map selfDescribedDsData = findInData(data, datasetUri)
String dsId = null
if (selfDescribedDsData != null) {
System.err.println("Using self-described dataset description")
log.info("Using self-described dataset description")
setDatasetInfo(datasetUri, data)
} else if (givenDsData != null) {
System.err.println("Using given dataset description")
log.info("Using given dataset description")
setDatasetInfo(datasetUri, givenDsData)
dsRecord = completeRecord(givenDsData, JsonLd.SYSTEM_RECORD_TYPE)
createOrUpdateDocument(dsRecord)
dsId = dsRecord.getShortId()
} else if (useExistingDatasetDescription) {
System.err.println("Using existing dataset description")
log.info("Using existing dataset description")
lookupDatasetInfo(datasetUri)
}
return dsId
Expand All @@ -272,7 +285,7 @@ class DatasetImporter {
throw new RuntimeException("Provided dataset ${givenData[ID]} does not match: ${datasetUri}")
}
dsInfo = new DatasetInfo(dsData)
System.err.println("Using new dataset: ${dsInfo.uri}")
log.info("Using new dataset: ${dsInfo.uri}")
}

protected void lookupDatasetInfo(String datasetUri) {
Expand All @@ -283,7 +296,7 @@ class DatasetImporter {
Map datasetData = ((List) datasetRecord.data[GRAPH])[1]
assert datasetData[ID] == datasetUri
dsInfo = new DatasetInfo(datasetData)
System.err.println("Using already defined dataset: ${dsInfo.uri}")
log.info("Using already defined dataset: ${dsInfo.uri}")
}

protected Document completeRecord(Map data, String recordType, boolean remap = false) {
Expand Down Expand Up @@ -372,9 +385,14 @@ class DatasetImporter {
}
}

/**
* This method is <em>only</em> intended for initial whelk bootstrapping,
* to initially load systemContext (contextDocData).
*/
private static Map loadSelfCompactedTurtle(InputStream ins) {
// Assuming that the Turtle *shape* follows a hard-coded system context!
Map data = (Map) TrigToJsonLdParser.parse(ins)
log.warn("Using loadSelfCompactedTurtle should only happen during setup of fresh whelk installation")
Map data = (Map) TrigToJsonLdParser.parseRaw(ins)
if (CONTEXT in data) {
Map ctx = [:]
ctx.putAll((Map) data[CONTEXT])
Expand All @@ -384,18 +402,27 @@ class DatasetImporter {
ctx['xsd'] = XSD_NS
// Assumes VOCAB + created in source actually means this!
ctx['created'] = [(TYPE): 'xsd:dateTime']
data = (Map) TrigToJsonLdParser.compact(data, [(CONTEXT): ctx])
data = (Map) JsonLdShapes.reCompact(data, [(CONTEXT): ctx])
}
return data
}

private Map loadTurtleAsSystemShaped(InputStream ins) {
assert contextDocData
Map data = TrigToJsonLdParser.parse(ins)
Map data = TrigToJsonLdParser.parseRaw(ins)
if (checkedSystemShaped(data, whelk.jsonld.vocabId)) {
return (Map) JsonLdShapes.reCompact(data, contextDocData)
} else {
log.info("Applying target vocabulary map")
return (Map) getTvm().applyTargetVocabularyMap(whelk.systemContextUri, contextDocData, data)
}
}

private static boolean checkedSystemShaped(Map data, String vocabId) {
if (data[CONTEXT] instanceof Map) {
Map ctx = (Map) data[CONTEXT]
int expectedSize = 0
if (ctx[VOCAB] == whelk.jsonld.vocabId) {
if (ctx[VOCAB] == vocabId) {
expectedSize++
if (ctx.containsKey(BASE)) {
expectedSize++
Expand All @@ -405,14 +432,15 @@ class DatasetImporter {
}
}
if (ctx.size() == expectedSize) {
// Forces plain string uri values to be taken as datatyped
// TODO: remove this hack when definitions consistently use `:uri ""^^xsd:anyURI`!
// Force plain string uri value to be expanded as datatyped:
if ('uri' !in ctx) {
ctx['uri'] = [(TYPE): 'xsd:anyURI']
}
return (Map) TrigToJsonLdParser.compact(data, contextDocData)
return true
}
}
return (Map) getTvm().applyTargetVocabularyMap(whelk.systemContextUri, contextDocData, data)
return false
}

private Document getDocByMainEntityId(String id) {
Expand Down Expand Up @@ -477,7 +505,7 @@ class DatasetImporter {
} else {
deletedCount++
if (deletedCount % 50 == 0) {
System.err.println("Cleaning up: " + deletedCount + " records deleted (they are no longer in the dataset).")
log.info("Cleaning up: " + deletedCount + " records deleted (they are no longer in the dataset).")
}
}
}
Expand Down
57 changes: 57 additions & 0 deletions importers/src/main/groovy/whelk/importer/ImporterMain.groovy
Original file line number Diff line number Diff line change
Expand Up @@ -3,18 +3,29 @@ package whelk.importer
import java.lang.annotation.*
import java.util.concurrent.ExecutorService
import java.util.zip.GZIPOutputStream

import java.nio.file.FileVisitResult
import java.nio.file.Files
import java.nio.file.Path
import java.nio.file.Paths
import java.nio.file.SimpleFileVisitor
import java.nio.file.attribute.BasicFileAttributes

import groovy.cli.picocli.CliBuilder
import groovy.util.logging.Slf4j as Log

import org.apache.commons.io.output.CountingOutputStream
import org.apache.commons.io.FilenameUtils

import whelk.Document
import whelk.Whelk
import whelk.component.PostgreSQLComponent
import whelk.converter.JsonLdToTrigSerializer
import whelk.converter.RdfReader
import whelk.filter.LinkFinder
import whelk.reindexer.CardRefresher
import whelk.reindexer.ElasticReindexer
import whelk.util.Jackson
import whelk.util.PropertyLoader

@Log
Expand All @@ -30,6 +41,52 @@ class ImporterMain {
return Whelk.createLoadedSearchWhelk(props)
}

@Command(args='SOURCE_URL [SUFFIXES]')
void rdfDirToJsonLines(String sourceDir, String suffixes='rdf,ttl,jsonld') {
var whelk = Whelk.createLoadedCoreWhelk(props)

var systemContextUri = whelk.systemContextUri
var baseUri = whelk.baseUri.toString()
System.err.println "Whelk system context URI: $systemContextUri; base URI: $baseUri"

var context = whelk.storage.loadDocumentByMainId(systemContextUri).data

var sourceDirPath = Paths.get(sourceDir)
var matcher = sourceDirPath.getFileSystem().getPathMatcher("glob:**/*.{$suffixes}")

var printStream = System.out

Files.walkFileTree(sourceDirPath, new SimpleFileVisitor<Path>() {
@Override
FileVisitResult visitFile(Path path, BasicFileAttributes attrs) {
if (!matcher.matches(path)) {
return FileVisitResult.CONTINUE
}

var rdfSourcePath = path.toString()
var data = new File(rdfSourcePath).withInputStream {
RdfReader.readRdf(it, rdfSourcePath, context, systemContextUri, baseUri)
}

if ('@graph' in data) {
// Ensure LDDB-expected order of things:
Map record = data['@graph'].find { it -> it['@type'] == 'Record' }
if (record) {
// assumes correctly structured record
def mainId = record['mainEntity']['@id']
Map mainEntity = data['@graph'].find { it -> it['@id'] == mainId }
List rest = data['@graph'].findAll { it -> !it.is(record) && !it.is(mainEntity) }
data = ['@graph': [record, mainEntity] + rest]
}
}

printStream.println(Jackson.mapper.writeValueAsString(data))

return FileVisitResult.CONTINUE
}
})
}

@Command(args='SOURCE_URL DATASET_URI [DATASET_DESCRIPTION_FILE]',
flags='--skip-index --replace-main-ids --force-delete --skip-dependers --allow-id-removal')
void dataset(Map flags, String sourceUrl, String datasetUri, String datasetDescPath=null) {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
package whelk.importer

import spock.lang.Specification

import whelk.converter.TrigToJsonLdParser

class DatasetImporterSpec extends Specification {

def "load self compacted turtle"() {
given:
var s = """
prefix : <https://id.kb.se/vocab/>
<x> a :Record .
"""
var data = DatasetImporter.loadSelfCompactedTurtle(new ByteArrayInputStream(s.getBytes('utf-8')))

expect:
data == [
"@id": "x",
"@type": "Record"
]
}

def "check system-shaped turtle"() {
when:
var s = """
prefix : <https://id.kb.se/vocab/>
prefix xsd: <http://www.w3.org/2001/XMLSchema#>

<ds/1> a :Dataset ;
:uri "ds/1" ; # NOTE: plain string form (notice difference in JSON-LD)
:created "2026-10-01T01:00:00Z"^^xsd:dateTime .
<ds/2> a :Dataset ;
:uri "ds/2"^^xsd:anyURI ;
:created "2026-10-01T02:00:00Z"^^xsd:dateTime .
"""
var data = TrigToJsonLdParser.parse(new ByteArrayInputStream(s.getBytes('utf-8')))

then:
data == [
"@context": [
"@vocab": "https://id.kb.se/vocab/",
"xsd": "http://www.w3.org/2001/XMLSchema#"
],
"@graph": [
[
"@id": "ds/1",
"@type": "Dataset",
"uri": "ds/1" ,
"created": ["@type": "xsd:dateTime", "@value": "2026-10-01T01:00:00Z"]
],
[
"@id": "ds/2",
"@type": "Dataset",
"uri": ["@type": "xsd:anyURI", "@value": "ds/2"],
"created": ["@type": "xsd:dateTime", "@value": "2026-10-01T02:00:00Z"]
]
]
]

and:
DatasetImporter.checkedSystemShaped(data, "https://id.kb.se/vocab/")

and:
data["@context"] == [
"@vocab": "https://id.kb.se/vocab/",
"xsd": "http://www.w3.org/2001/XMLSchema#",
"uri": ["@type": "xsd:anyURI"]
]
}

}
Loading
Loading