-
Notifications
You must be signed in to change notification settings - Fork 8.2k
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
Showing
5 changed files
with
288 additions
and
7 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
198 changes: 198 additions & 0 deletions
198
x-pack/legacy/plugins/ml/server/models/file_data_visualizer/import_data.ts
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,198 @@ | ||
/* | ||
* Copyright Elasticsearch B.V. and/or licensed to Elasticsearch B.V. under one | ||
* or more contributor license agreements. Licensed under the Elastic License; | ||
* you may not use this file except in compliance with the Elastic License. | ||
*/ | ||
|
||
import { RequestHandlerContext } from 'kibana/server'; | ||
import { INDEX_META_DATA_CREATED_BY } from '../../../common/constants/file_datavisualizer'; | ||
import { InputData } from './file_data_visualizer'; | ||
|
||
export interface Settings { | ||
pipeline?: string; | ||
index: string; | ||
body: any[]; | ||
[key: string]: any; | ||
} | ||
|
||
export interface Mappings { | ||
[key: string]: any; | ||
} | ||
|
||
export interface InjectPipeline { | ||
id: string; | ||
pipeline: any; | ||
} | ||
|
||
interface Failure { | ||
item: number; | ||
reason: string; | ||
doc: any; | ||
} | ||
|
||
export function importDataProvider(context: RequestHandlerContext) { | ||
const callAsCurrentUser = context.ml!.mlClient.callAsCurrentUser; | ||
|
||
async function importData( | ||
id: string, | ||
index: string, | ||
settings: Settings, | ||
mappings: Mappings, | ||
ingestPipeline: InjectPipeline, | ||
data: InputData | ||
) { | ||
let createdIndex; | ||
let createdPipelineId; | ||
const docCount = data.length; | ||
|
||
try { | ||
const { id: pipelineId, pipeline } = ingestPipeline; | ||
|
||
if (id === undefined) { | ||
// first chunk of data, create the index and id to return | ||
id = generateId(); | ||
|
||
await createIndex(index, settings, mappings); | ||
createdIndex = index; | ||
|
||
// create the pipeline if one has been supplied | ||
if (pipelineId !== undefined) { | ||
const success = await createPipeline(pipelineId, pipeline); | ||
if (success.acknowledged !== true) { | ||
throw success; | ||
} | ||
} | ||
createdPipelineId = pipelineId; | ||
} else { | ||
createdIndex = index; | ||
createdPipelineId = pipelineId; | ||
} | ||
|
||
let failures: Failure[] = []; | ||
if (data.length) { | ||
const resp = await indexData(index, createdPipelineId, data); | ||
if (resp.success === false) { | ||
if (resp.ingestError) { | ||
// all docs failed, abort | ||
throw resp; | ||
} else { | ||
// some docs failed. | ||
// still report success but with a list of failures | ||
failures = resp.failures || []; | ||
} | ||
} | ||
} | ||
|
||
return { | ||
success: true, | ||
id, | ||
index: createdIndex, | ||
pipelineId: createdPipelineId, | ||
docCount, | ||
failures, | ||
}; | ||
} catch (error) { | ||
return { | ||
success: false, | ||
id, | ||
index: createdIndex, | ||
pipelineId: createdPipelineId, | ||
error: error.error !== undefined ? error.error : error, | ||
docCount, | ||
ingestError: error.ingestError, | ||
failures: error.failures || [], | ||
}; | ||
} | ||
} | ||
|
||
async function createIndex(index: string, settings: Settings, mappings: Mappings) { | ||
const body: { mappings: Mappings; settings?: Settings } = { | ||
mappings: { | ||
_meta: { | ||
created_by: INDEX_META_DATA_CREATED_BY, | ||
}, | ||
properties: mappings, | ||
}, | ||
}; | ||
|
||
if (settings && Object.keys(settings).length) { | ||
body.settings = settings; | ||
} | ||
|
||
await callAsCurrentUser('indices.create', { index, body }); | ||
} | ||
|
||
async function indexData(index: string, pipelineId: string, data: InputData) { | ||
try { | ||
const body = []; | ||
for (let i = 0; i < data.length; i++) { | ||
body.push({ index: {} }); | ||
body.push(data[i]); | ||
} | ||
|
||
const settings: Settings = { index, body }; | ||
if (pipelineId !== undefined) { | ||
settings.pipeline = pipelineId; | ||
} | ||
|
||
const resp = await callAsCurrentUser('bulk', settings); | ||
if (resp.errors) { | ||
throw resp; | ||
} else { | ||
return { | ||
success: true, | ||
docs: data.length, | ||
failures: [], | ||
}; | ||
} | ||
} catch (error) { | ||
let failures: Failure[] = []; | ||
let ingestError = false; | ||
if (error.errors !== undefined && Array.isArray(error.items)) { | ||
// an expected error where some or all of the bulk request | ||
// docs have failed to be ingested. | ||
failures = getFailures(error.items, data); | ||
} else { | ||
// some other error has happened. | ||
ingestError = true; | ||
} | ||
|
||
return { | ||
success: false, | ||
error, | ||
docCount: data.length, | ||
failures, | ||
ingestError, | ||
}; | ||
} | ||
} | ||
|
||
async function createPipeline(id: string, pipeline: any) { | ||
return await callAsCurrentUser('ingest.putPipeline', { id, body: pipeline }); | ||
} | ||
|
||
function getFailures(items: any[], data: InputData): Failure[] { | ||
const failures = []; | ||
for (let i = 0; i < items.length; i++) { | ||
const item = items[i]; | ||
if (item.index && item.index.error) { | ||
failures.push({ | ||
item: i, | ||
reason: item.index.error.reason, | ||
doc: data[i], | ||
}); | ||
} | ||
} | ||
return failures; | ||
} | ||
|
||
return { | ||
importData, | ||
}; | ||
} | ||
|
||
function generateId() { | ||
return Math.random() | ||
.toString(36) | ||
.substr(2, 9); | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters