|
16 | 16 | * specific language governing permissions and limitations |
17 | 17 | * under the License. |
18 | 18 | */ |
19 | | -package org.apache.polaris.service.lineage; |
| 19 | +package org.apache.polaris.extensions.lineage; |
20 | 20 |
|
21 | 21 | import jakarta.enterprise.context.RequestScoped; |
22 | 22 | import jakarta.inject.Inject; |
23 | 23 | import java.time.Instant; |
24 | 24 | import org.apache.polaris.core.config.FeatureConfiguration; |
25 | 25 | import org.apache.polaris.core.context.CallContext; |
26 | 26 | import org.apache.polaris.core.context.RealmContext; |
27 | | -import org.apache.polaris.core.lineage.LineageGraph; |
28 | | -import org.apache.polaris.core.lineage.LineageIngestRequest; |
29 | | -import org.apache.polaris.core.lineage.LineagePersistence; |
30 | | -import org.apache.polaris.core.lineage.LineageQueryRequest; |
31 | | -import org.apache.polaris.core.lineage.LineageService; |
32 | 27 |
|
33 | 28 | @RequestScoped |
34 | | -public class DefaultLineageService implements LineageService { |
| 29 | +public class DefaultPolarisLineageHandler implements PolarisLineageHandler { |
35 | 30 | private final CallContext callContext; |
36 | 31 | private final LineageConfiguration configuration; |
37 | | - private final LineagePersistence persistence; |
| 32 | + private final LineageStoreManager storeManager; |
38 | 33 |
|
39 | 34 | @Inject |
40 | | - public DefaultLineageService( |
41 | | - CallContext callContext, LineageConfiguration configuration, LineagePersistence persistence) { |
| 35 | + public DefaultPolarisLineageHandler( |
| 36 | + CallContext callContext, |
| 37 | + LineageConfiguration configuration, |
| 38 | + LineageStoreManager storeManager) { |
42 | 39 | this.callContext = callContext; |
43 | 40 | this.configuration = configuration; |
44 | | - this.persistence = persistence; |
| 41 | + this.storeManager = storeManager; |
45 | 42 | } |
46 | 43 |
|
47 | 44 | @Override |
48 | 45 | public void ingest(LineageIngestRequest request) { |
49 | 46 | ensureEnabled(); |
50 | 47 | RealmContext realmContext = callContext.getRealmContext(); |
51 | 48 | Instant lastEventAt = request.eventTime().orElseGet(Instant::now); |
52 | | - persistence.upsertDatasets(realmContext, request.datasets()); |
53 | | - persistence.replaceDatasetEdges( |
| 49 | + storeManager.upsertDatasets(realmContext, request.datasets()); |
| 50 | + storeManager.replaceDatasetEdges( |
54 | 51 | realmContext, request.targetDatasets(), request.edges(), lastEventAt); |
55 | | - persistence.upsertColumnEdges(realmContext, request.columnEdges(), lastEventAt); |
| 52 | + storeManager.upsertColumnEdges(realmContext, request.columnEdges(), lastEventAt); |
56 | 53 | } |
57 | 54 |
|
58 | 55 | @Override |
59 | 56 | public LineageGraph query(LineageQueryRequest request) { |
60 | 57 | ensureEnabled(); |
61 | | - return persistence.loadLineage(callContext.getRealmContext(), request); |
| 58 | + return storeManager.loadLineage(callContext.getRealmContext(), request); |
62 | 59 | } |
63 | 60 |
|
64 | 61 | private void ensureEnabled() { |
|
0 commit comments